From c7281f27e0d62a7204acbfaaca1b9226a3cf371f Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Sat, 27 Jul 2019 00:18:13 -0700 Subject: [PATCH 1/2] Change back to normal number writings as not supported by python under 3.6 --- .../azure-eventhubs/azure/eventhub/aio/consumer_async.py | 2 +- .../azure-eventhubs/azure/eventhub/aio/producer_async.py | 2 +- sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py | 2 +- sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py index 286b62f0a05c..dac2d0c0fa61 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/consumer_async.py @@ -193,7 +193,7 @@ async def receive(self, **kwargs): max_batch_size = min(self.client.config.max_batch_size, self.prefetch) if max_batch_size is None else max_batch_size timeout = self.client.config.receive_timeout if timeout is None else timeout if not timeout: - timeout = 100_000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout + timeout = 100000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout data_batch = [] start_time = time.time() diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py index bb16b8ee45bb..e326aef0a115 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/producer_async.py @@ -113,7 +113,7 @@ async def _open(self, timeout_time=None): async def _send_event_data(self, timeout=None): timeout = timeout or self.client.config.send_timeout if not timeout: - timeout = 100_000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout + timeout = 100000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout start_time = time.time() timeout_time = start_time + timeout max_retries = self.client.config.max_retries diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py index dd83b3848748..e59c440c7c88 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/consumer.py @@ -185,7 +185,7 @@ def receive(self, **kwargs): max_batch_size = min(self.client.config.max_batch_size, self.prefetch) if max_batch_size is None else max_batch_size timeout = self.client.config.receive_timeout if timeout is None else timeout if not timeout: - timeout = 100_000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout + timeout = 100000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout data_batch = [] # type: List[EventData] start_time = time.time() diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py index 16803dace84e..da2a9ee95368 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/producer.py @@ -121,7 +121,7 @@ def _open(self, timeout_time=None): def _send_event_data(self, timeout=None): timeout = timeout or self.client.config.send_timeout if not timeout: - timeout = 100_000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout + timeout = 100000 # timeout None or 0 mean no timeout. 100000 seconds is equivalent to no timeout start_time = time.time() timeout_time = start_time + timeout max_retries = self.client.config.max_retries From 51e53beaddd5e1d492ed72f3ecf7d9d3292854a2 Mon Sep 17 00:00:00 2001 From: Yunhao Ling Date: Sun, 28 Jul 2019 20:31:26 -0700 Subject: [PATCH 2/2] small fix --- .../eventhub/aio/_connection_manager_async.py | 2 +- .../azure-eventhubs/azure/eventhub/error.py | 3 +- .../tests/asynctests/test_auth_async.py | 4 +- .../tests/asynctests/test_negative_async.py | 38 +++++++++---------- .../azure-eventhubs/tests/test_auth.py | 4 +- .../azure-eventhubs/tests/test_negative.py | 22 +++++++---- 6 files changed, 38 insertions(+), 35 deletions(-) diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_connection_manager_async.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_connection_manager_async.py index 989f3416d582..618359192ffe 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_connection_manager_async.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/aio/_connection_manager_async.py @@ -70,7 +70,7 @@ async def get_connection(self, host, auth): async def close_connection(self): pass - def reset_connection_if_broken(self): + async def reset_connection_if_broken(self): pass diff --git a/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py b/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py index 0fb6933e3015..cbe6a8a04946 100644 --- a/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py +++ b/sdk/eventhub/azure-eventhubs/azure/eventhub/error.py @@ -57,8 +57,7 @@ class EventHubError(Exception): :vartype details: dict[str, str] """ - def __init__(self, message, **kwargs): - details = kwargs.get("details", None) + def __init__(self, message, details=None): self.error = None self.message = message self.details = details diff --git a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_auth_async.py b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_auth_async.py index 0a1ef23e1dea..d0a0d0912868 100644 --- a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_auth_async.py +++ b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_auth_async.py @@ -31,7 +31,7 @@ async def test_client_secret_credential_async(aad_credential, live_eventhub): async with receiver: - received = await receiver.receive(timeout=1) + received = await receiver.receive(timeout=3) assert len(received) == 0 async with sender: @@ -40,7 +40,7 @@ async def test_client_secret_credential_async(aad_credential, live_eventhub): await asyncio.sleep(1) - received = await receiver.receive(timeout=1) + received = await receiver.receive(timeout=3) assert len(received) == 1 assert list(received[0].body)[0] == 'A single message'.encode('utf-8') diff --git a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_negative_async.py b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_negative_async.py index 27deb42c3aeb..4504d8558298 100644 --- a/sdk/eventhub/azure-eventhubs/tests/asynctests/test_negative_async.py +++ b/sdk/eventhub/azure-eventhubs/tests/asynctests/test_negative_async.py @@ -29,7 +29,7 @@ async def test_send_with_invalid_hostname_async(invalid_hostname, connstr_receiv client = EventHubClient.from_connection_string(invalid_hostname, network_tracing=False) sender = client.create_producer() with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -38,7 +38,7 @@ async def test_receive_with_invalid_hostname_async(invalid_hostname): client = EventHubClient.from_connection_string(invalid_hostname, network_tracing=False) sender = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -48,7 +48,7 @@ async def test_send_with_invalid_key_async(invalid_key, connstr_receivers): client = EventHubClient.from_connection_string(invalid_key, network_tracing=False) sender = client.create_producer() with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -57,7 +57,7 @@ async def test_receive_with_invalid_key_async(invalid_key): client = EventHubClient.from_connection_string(invalid_key, network_tracing=False) sender = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -67,7 +67,7 @@ async def test_send_with_invalid_policy_async(invalid_policy, connstr_receivers) client = EventHubClient.from_connection_string(invalid_policy, network_tracing=False) sender = client.create_producer() with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -76,7 +76,7 @@ async def test_receive_with_invalid_policy_async(invalid_policy): client = EventHubClient.from_connection_string(invalid_policy, network_tracing=False) sender = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -88,7 +88,7 @@ async def test_send_partition_key_with_partition_async(connection_str): try: data = EventData(b"Data") with pytest.raises(ValueError): - await sender.send(data) + await sender.send(EventData("test data")) finally: await sender.close() @@ -99,7 +99,7 @@ async def test_non_existing_entity_sender_async(connection_str): client = EventHubClient.from_connection_string(connection_str, event_hub_path="nemo", network_tracing=False) sender = client.create_producer(partition_id="1") with pytest.raises(AuthenticationError): - await sender._open() + await sender.send(EventData("test data")) @pytest.mark.liveTest @@ -108,35 +108,31 @@ async def test_non_existing_entity_receiver_async(connection_str): client = EventHubClient.from_connection_string(connection_str, event_hub_path="nemo", network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - await receiver._open() + await receiver.receive(timeout=5) @pytest.mark.liveTest @pytest.mark.asyncio async def test_receive_from_invalid_partitions_async(connection_str): - partitions = ["XYZ", "-1", "1000", "-" ] + partitions = ["XYZ", "-1", "1000", "-"] for p in partitions: client = EventHubClient.from_connection_string(connection_str, network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id=p, event_position=EventPosition("-1")) - try: - with pytest.raises(ConnectError): - await receiver.receive(timeout=10) - finally: - await receiver.close() + with pytest.raises(ConnectError): + await receiver.receive(timeout=10) + await receiver.close() @pytest.mark.liveTest @pytest.mark.asyncio async def test_send_to_invalid_partitions_async(connection_str): - partitions = ["XYZ", "-1", "1000", "-" ] + partitions = ["XYZ", "-1", "1000", "-"] for p in partitions: client = EventHubClient.from_connection_string(connection_str, network_tracing=False) sender = client.create_producer(partition_id=p) - try: - with pytest.raises(ConnectError): - await sender._open() - finally: - await sender.close() + with pytest.raises(ConnectError): + await sender.send(EventData("test data")) + await sender.close() @pytest.mark.liveTest diff --git a/sdk/eventhub/azure-eventhubs/tests/test_auth.py b/sdk/eventhub/azure-eventhubs/tests/test_auth.py index d5871971a5b4..0f0c78794884 100644 --- a/sdk/eventhub/azure-eventhubs/tests/test_auth.py +++ b/sdk/eventhub/azure-eventhubs/tests/test_auth.py @@ -26,7 +26,7 @@ def test_client_secret_credential(aad_credential, live_eventhub): receiver = client.create_consumer(consumer_group="$default", partition_id='0', event_position=EventPosition("@latest")) with receiver: - received = receiver.receive(timeout=1) + received = receiver.receive(timeout=3) assert len(received) == 0 with sender: @@ -34,7 +34,7 @@ def test_client_secret_credential(aad_credential, live_eventhub): sender.send(event) time.sleep(1) - received = receiver.receive(timeout=1) + received = receiver.receive(timeout=3) assert len(received) == 1 assert list(received[0].body)[0] == 'A single message'.encode('utf-8') diff --git a/sdk/eventhub/azure-eventhubs/tests/test_negative.py b/sdk/eventhub/azure-eventhubs/tests/test_negative.py index ac19a01f76c9..4749df940d9c 100644 --- a/sdk/eventhub/azure-eventhubs/tests/test_negative.py +++ b/sdk/eventhub/azure-eventhubs/tests/test_negative.py @@ -34,7 +34,8 @@ def test_receive_with_invalid_hostname_sync(invalid_hostname): client = EventHubClient.from_connection_string(invalid_hostname, network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - receiver.receive(timeout=3) + receiver.receive(timeout=5) + receiver.close() @pytest.mark.liveTest @@ -44,14 +45,16 @@ def test_send_with_invalid_key(invalid_key, connstr_receivers): sender = client.create_producer() with pytest.raises(AuthenticationError): sender.send(EventData("test data")) - + sender.close() @pytest.mark.liveTest def test_receive_with_invalid_key_sync(invalid_key): client = EventHubClient.from_connection_string(invalid_key, network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) + with pytest.raises(AuthenticationError): - receiver.receive(timeout=3) + receiver.receive(timeout=10) + receiver.close() @pytest.mark.liveTest @@ -61,6 +64,7 @@ def test_send_with_invalid_policy(invalid_policy, connstr_receivers): sender = client.create_producer() with pytest.raises(AuthenticationError): sender.send(EventData("test data")) + sender.close() @pytest.mark.liveTest @@ -68,7 +72,8 @@ def test_receive_with_invalid_policy_sync(invalid_policy): client = EventHubClient.from_connection_string(invalid_policy, network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) with pytest.raises(AuthenticationError): - receiver.receive(timeout=3) + receiver.receive(timeout=5) + receiver.close() @pytest.mark.liveTest @@ -97,13 +102,16 @@ def test_non_existing_entity_sender(connection_str): def test_non_existing_entity_receiver(connection_str): client = EventHubClient.from_connection_string(connection_str, event_hub_path="nemo", network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id="0", event_position=EventPosition("-1")) + with pytest.raises(AuthenticationError): - receiver.receive(timeout=3) + receiver.receive(timeout=5) + receiver.close() + @pytest.mark.liveTest def test_receive_from_invalid_partitions_sync(connection_str): - partitions = ["XYZ", "-1", "1000", "-" ] + partitions = ["XYZ", "-1", "1000", "-"] for p in partitions: client = EventHubClient.from_connection_string(connection_str, network_tracing=False) receiver = client.create_consumer(consumer_group="$default", partition_id=p, event_position=EventPosition("-1")) @@ -116,7 +124,7 @@ def test_receive_from_invalid_partitions_sync(connection_str): @pytest.mark.liveTest def test_send_to_invalid_partitions(connection_str): - partitions = ["XYZ", "-1", "1000", "-" ] + partitions = ["XYZ", "-1", "1000", "-"] for p in partitions: client = EventHubClient.from_connection_string(connection_str, network_tracing=False) sender = client.create_producer(partition_id=p)