From bd934d983c118555ae84d8a94728ec252108f8b1 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:08:04 -0500 Subject: [PATCH 01/14] mess around this unit tests for aiokafka statuses --- faust/transport/drivers/aiokafka.py | 4 +- tests/unit/transport/drivers/test_aiokafka.py | 46 ++++++++++++++++--- 2 files changed, 42 insertions(+), 8 deletions(-) diff --git a/faust/transport/drivers/aiokafka.py b/faust/transport/drivers/aiokafka.py index f69c532db..eaa16d2e9 100644 --- a/faust/transport/drivers/aiokafka.py +++ b/faust/transport/drivers/aiokafka.py @@ -26,7 +26,7 @@ import aiokafka import aiokafka.abc import opentracing -from aiokafka import TopicPartition +from aiokafka import TopicPartition, AIOKafkaConsumer from aiokafka.consumer.group_coordinator import OffsetCommitRequest from aiokafka.coordinator.assignors.roundrobin import RoundRobinPartitionAssignor from aiokafka.errors import ( @@ -830,7 +830,7 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: Returns :const:`True` if any error was logged. """ - consumer = self._ensure_consumer() + consumer: AIOKafkaConsumer = self._ensure_consumer() secs_since_started = now - self.time_started aiotp = TopicPartition(tp.topic, tp.partition) assignment = consumer._fetcher._subscriptions.subscription.assignment diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 63a55d5d2..fd67b5d50 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -9,6 +9,7 @@ import pytest from aiokafka.errors import CommitFailedError, IllegalStateError, KafkaError from aiokafka.structs import OffsetAndMetadata, TopicPartition +from aiokafka.consumer.subscription_state import TopicPartitionState from mode.utils import text from mode.utils.futures import done_future from opentracing.ext import tags @@ -234,7 +235,7 @@ def mock_record( serialized_value_size=40, **kwargs, ): - return Mock( + return MagicMock( name="record", topic=topic, partition=partition, @@ -503,29 +504,62 @@ def test_timed_out(self, *, cthread, now, tp, logger): ) -@pytest.mark.skip("Needs fixing") class Test_VEP_no_highwater_since_start(Test_verify_event_path_base): highwater = None - def test_no_monitor(self, *, app, cthread, now, tp, logger): + def test_no_monitor(self, *, app, cthread, now, tp, logger, _consumer): self._set_last_request(now - 10.0) self._set_last_response(now - 5.0) self._set_started(now) app.monitor = None + _consumer.assignment.return_value = {tp} + assignment = cthread.assignment() + + assert assignment == {tp} + _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( + assignment=assignment, + timestamp=now, + tp_fetch_request_timeout_secs=10.0, + ) + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() - def test_just_started(self, *, cthread, now, tp, logger): + def test_just_started(self, *, cthread, now, tp, logger, _consumer): self._set_last_request(now - 10.0) self._set_last_response(now - 5.0) self._set_started(now) + _consumer.assignment.return_value = {tp} + assignment = cthread.assignment() + + assert assignment == {tp} + _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( + assignment=assignment, + timestamp=now, + tp_fetch_request_timeout_secs=10.0, + ) + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() - def test_timed_out(self, *, cthread, now, tp, logger): + def test_timed_out(self, *, cthread, now, tp, logger, _consumer): self._set_last_request(now - 10.0) self._set_last_response(now - 5.0) self._set_started(now - cthread.tp_stream_timeout_secs * 2) + _consumer.assignment.return_value = {tp} + assignment = cthread.assignment() + + assert assignment == {tp} + _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( + assignment=assignment, + timestamp=now, + highwater=None, + tp_stream_timeout_secs=cthread.tp_stream_timeout_secs, + tp_fetch_request_timeout_secs=cthread.tp_fetch_request_timeout_secs, + ) + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + assert cthread.verify_event_path(now, tp) is None logger.error.assert_called_with( mod.SLOW_PROCESSING_NO_HIGHWATER_SINCE_START, @@ -636,7 +670,7 @@ def test_inbound_timed_out(self, *, app, cthread, now, tp, logger): ) -@pytest.mark.skip("Needs fixing") +# @pytest.mark.skip("Needs fixing") class Test_VEP_no_commit(Test_verify_event_path_base): highwater = 20 committed_offset = 10 From 75e8f875b64ebefc1aefb49c43563b34e23c2113 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:13:35 -0500 Subject: [PATCH 02/14] apparently highwater had a bug??? --- faust/transport/drivers/aiokafka.py | 4 ++-- tests/unit/transport/drivers/test_aiokafka.py | 17 ++++++++++++----- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/faust/transport/drivers/aiokafka.py b/faust/transport/drivers/aiokafka.py index eaa16d2e9..57ff3ca0f 100644 --- a/faust/transport/drivers/aiokafka.py +++ b/faust/transport/drivers/aiokafka.py @@ -26,7 +26,7 @@ import aiokafka import aiokafka.abc import opentracing -from aiokafka import TopicPartition, AIOKafkaConsumer +from aiokafka import AIOKafkaConsumer, TopicPartition from aiokafka.consumer.group_coordinator import OffsetCommitRequest from aiokafka.coordinator.assignors.roundrobin import RoundRobinPartitionAssignor from aiokafka.errors import ( @@ -840,7 +840,7 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: poll_at = None aiotp_state = assignment.state_value(aiotp) if aiotp_state and aiotp_state.timestamp: - poll_at = aiotp_state.timestamp / 1000 + poll_at = aiotp_state.timestamp if poll_at is None: if secs_since_started >= self.tp_fetch_request_timeout_secs: # NO FETCH REQUEST SENT AT ALL SINCE WORKER START diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index fd67b5d50..d74bac64a 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -7,9 +7,9 @@ import aiokafka import opentracing import pytest +from aiokafka.consumer.subscription_state import TopicPartitionState from aiokafka.errors import CommitFailedError, IllegalStateError, KafkaError from aiokafka.structs import OffsetAndMetadata, TopicPartition -from aiokafka.consumer.subscription_state import TopicPartitionState from mode.utils import text from mode.utils.futures import done_future from opentracing.ext import tags @@ -521,7 +521,9 @@ def test_no_monitor(self, *, app, cthread, now, tp, logger, _consumer): timestamp=now, tp_fetch_request_timeout_secs=10.0, ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( + now + ) assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() @@ -538,7 +540,9 @@ def test_just_started(self, *, cthread, now, tp, logger, _consumer): timestamp=now, tp_fetch_request_timeout_secs=10.0, ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( + now + ) assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() @@ -558,7 +562,9 @@ def test_timed_out(self, *, cthread, now, tp, logger, _consumer): tp_stream_timeout_secs=cthread.tp_stream_timeout_secs, tp_fetch_request_timeout_secs=cthread.tp_fetch_request_timeout_secs, ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( + now + ) assert cthread.verify_event_path(now, tp) is None logger.error.assert_called_with( @@ -1322,7 +1328,8 @@ def assert_calls_thread(self, cthread, _consumer, method, *args, **kwargs): cthread.call_thread.assert_called_once_with(method, *args, **kwargs) -class MyPartitioner: ... +class MyPartitioner: + ... my_partitioner = MyPartitioner() From 492ea2e30749ec1b88344d01b6083905b68a4418 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:13:56 -0500 Subject: [PATCH 03/14] fix formatting --- tests/unit/transport/drivers/test_aiokafka.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index d74bac64a..50ccf0189 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -1328,8 +1328,7 @@ def assert_calls_thread(self, cthread, _consumer, method, *args, **kwargs): cthread.call_thread.assert_called_once_with(method, *args, **kwargs) -class MyPartitioner: - ... +class MyPartitioner: ... my_partitioner = MyPartitioner() From 4f1323b45a023f54445f24498e51e7e1523aa406 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:20:22 -0500 Subject: [PATCH 04/14] wow has this been broken this entire time??? --- tests/unit/transport/drivers/test_aiokafka.py | 42 ++++++------------- 1 file changed, 12 insertions(+), 30 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 50ccf0189..85185e9a6 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -282,8 +282,8 @@ def start_span(operation_name=None, **kwargs): return tracer @pytest.fixture() - def _consumer(self): - return Mock( + def _consumer(self, now, cthread, tp): + _consumer = Mock( name="AIOKafkaConsumer", autospec=aiokafka.AIOKafkaConsumer, start=AsyncMock(), @@ -294,6 +294,15 @@ def _consumer(self): _client=Mock(name="Client", close=AsyncMock()), _coordinator=Mock(name="Coordinator", close=AsyncMock()), ) + _consumer.assignment.return_value = {tp} + + _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( + assignment={tp}, + timestamp=now, + ) + # _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now + return _consumer + @pytest.fixture() def now(self): @@ -512,18 +521,6 @@ def test_no_monitor(self, *, app, cthread, now, tp, logger, _consumer): self._set_last_response(now - 5.0) self._set_started(now) app.monitor = None - _consumer.assignment.return_value = {tp} - assignment = cthread.assignment() - - assert assignment == {tp} - _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( - assignment=assignment, - timestamp=now, - tp_fetch_request_timeout_secs=10.0, - ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( - now - ) assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() @@ -531,19 +528,6 @@ def test_just_started(self, *, cthread, now, tp, logger, _consumer): self._set_last_request(now - 10.0) self._set_last_response(now - 5.0) self._set_started(now) - _consumer.assignment.return_value = {tp} - assignment = cthread.assignment() - - assert assignment == {tp} - _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( - assignment=assignment, - timestamp=now, - tp_fetch_request_timeout_secs=10.0, - ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( - now - ) - assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() @@ -574,7 +558,6 @@ def test_timed_out(self, *, cthread, now, tp, logger, _consumer): ) -@pytest.mark.skip("Needs fixing") class Test_VEP_stream_idle_no_highwater(Test_verify_event_path_base): highwater = 10 committed_offset = 10 @@ -587,7 +570,6 @@ def test_highwater_same_as_offset(self, *, cthread, now, tp, logger): logger.error.assert_not_called() -@pytest.mark.skip("Needs fixing") class Test_VEP_stream_idle_highwater_no_acks(Test_verify_event_path_base): acks_enabled = False @@ -599,7 +581,6 @@ def test_no_acks(self, *, cthread, now, tp, logger): logger.error.assert_not_called() -@pytest.mark.skip("Needs fixing") class Test_VEP_stream_idle_highwater_same_has_acks_everything_OK( Test_verify_event_path_base ): @@ -612,6 +593,7 @@ def test_main(self, *, cthread, now, tp, logger): self._set_last_request(now - 10.0) self._set_last_response(now - 5.0) self._set_started(now) + assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() From ce00568edf47449189304add583c37d0223f5b06 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:37:29 -0500 Subject: [PATCH 05/14] fix more tests --- tests/unit/transport/drivers/test_aiokafka.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 85185e9a6..d341562b1 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -7,10 +7,10 @@ import aiokafka import opentracing import pytest -from aiokafka.consumer.subscription_state import TopicPartitionState from aiokafka.errors import CommitFailedError, IllegalStateError, KafkaError from aiokafka.structs import OffsetAndMetadata, TopicPartition from mode.utils import text +from mode.utils.times import humanize_seconds_ago from mode.utils.futures import done_future from opentracing.ext import tags @@ -299,8 +299,9 @@ def _consumer(self, now, cthread, tp): _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( assignment={tp}, timestamp=now, + highwater=1, + position=0, ) - # _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = now return _consumer @@ -658,7 +659,6 @@ def test_inbound_timed_out(self, *, app, cthread, now, tp, logger): ) -# @pytest.mark.skip("Needs fixing") class Test_VEP_no_commit(Test_verify_event_path_base): highwater = 20 committed_offset = 10 @@ -686,13 +686,13 @@ def test_timed_out_since_start(self, *, app, cthread, now, tp, logger): expected_message = cthread._make_slow_processing_error( mod.SLOW_PROCESSING_NO_COMMIT_SINCE_START, [mod.SLOW_PROCESSING_CAUSE_COMMIT], + setting="broker_commit_livelock_soft_timeout", + current_value=app.conf.broker_commit_livelock_soft_timeout, ) logger.error.assert_called_once_with( expected_message, tp, - ANY, - setting="broker_commit_livelock_soft_timeout", - current_value=app.conf.broker_commit_livelock_soft_timeout, + humanize_seconds_ago(cthread.tp_commit_timeout_secs * 2), ) def test_timed_out_since_last(self, *, app, cthread, now, tp, logger): @@ -703,13 +703,13 @@ def test_timed_out_since_last(self, *, app, cthread, now, tp, logger): expected_message = cthread._make_slow_processing_error( mod.SLOW_PROCESSING_NO_RECENT_COMMIT, [mod.SLOW_PROCESSING_CAUSE_COMMIT], + setting="broker_commit_livelock_soft_timeout", + current_value=app.conf.broker_commit_livelock_soft_timeout, ) logger.error.assert_called_once_with( expected_message, tp, - ANY, - setting="broker_commit_livelock_soft_timeout", - current_value=app.conf.broker_commit_livelock_soft_timeout, + humanize_seconds_ago(now - cthread.tp_commit_timeout_secs * 4), ) def test_committing_fine(self, *, app, cthread, now, tp, logger): From cb6f4ef40684ce74263d1b0a09d91d121ef0e2a3 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 13:38:21 -0500 Subject: [PATCH 06/14] lint --- tests/unit/transport/drivers/test_aiokafka.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index d341562b1..87e59857b 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -10,8 +10,8 @@ from aiokafka.errors import CommitFailedError, IllegalStateError, KafkaError from aiokafka.structs import OffsetAndMetadata, TopicPartition from mode.utils import text -from mode.utils.times import humanize_seconds_ago from mode.utils.futures import done_future +from mode.utils.times import humanize_seconds_ago from opentracing.ext import tags import faust @@ -304,7 +304,6 @@ def _consumer(self, now, cthread, tp): ) return _consumer - @pytest.fixture() def now(self): return 1201230410 From a7cea114401a7360a65e18dccf96eaa71a62ed62 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 14:04:35 -0500 Subject: [PATCH 07/14] fix yet another test --- tests/unit/transport/drivers/test_aiokafka.py | 23 ++++++++++++++----- 1 file changed, 17 insertions(+), 6 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 87e59857b..c66deb579 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -1365,6 +1365,8 @@ def assert_new_producer( max_batch_size=16384, max_request_size=1000000, request_timeout_ms=1200000, + metadata_max_age_ms=300000, + connections_max_idle_ms=540000, security_protocol="PLAINTEXT", **kwargs, ): @@ -1372,19 +1374,21 @@ def assert_new_producer( p = producer._new_producer() assert p is AIOKafkaProducer.return_value AIOKafkaProducer.assert_called_once_with( - acks=acks, - api_version=api_version, bootstrap_servers=bootstrap_servers, client_id=client_id, - compression_type=compression_type, + acks=acks, linger_ms=linger_ms, max_batch_size=max_batch_size, max_request_size=max_request_size, - request_timeout_ms=request_timeout_ms, + compression_type=compression_type, security_protocol=security_protocol, loop=producer.loop, partitioner=producer.partitioner, transactional_id=None, + api_version=api_version, + metadata_max_age_ms=metadata_max_age_ms, + connections_max_idle_ms=connections_max_idle_ms, + request_timeout_ms=request_timeout_ms, **kwargs, ) @@ -1454,7 +1458,6 @@ def test__settings_extra(self, *, producer, app): app.in_transaction = False assert producer._settings_extra() == {} - @pytest.mark.skip("fix me") def test__new_producer(self, *, app): producer = Producer(app.transport) self.assert_new_producer(producer) @@ -1463,7 +1466,7 @@ def test__new_producer(self, *, app): "expected_args", [ pytest.param( - {"api_version": "0.10"}, + {"api_version": ""}, marks=pytest.mark.conf(producer_api_version="0.10"), ), pytest.param({"acks": -1}, marks=pytest.mark.conf(producer_acks="all")), @@ -1495,6 +1498,14 @@ def test__new_producer(self, *, app): {"request_timeout_ms": 1234134000}, marks=pytest.mark.conf(producer_request_timeout=1234134), ), + pytest.param( + {"metadata_max_age_ms": 300000}, + marks=pytest.mark.conf(metadata_max_age_ms=300000), + ), + pytest.param( + {"connections_max_idle_ms": 540000}, + marks=pytest.mark.conf(connections_max_idle_ms=540000), + ), pytest.param( { "security_protocol": "SASL_PLAINTEXT", From 56cad2ba84eb7d75a6f72623097d2cfa6d85224c Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 14:07:34 -0500 Subject: [PATCH 08/14] fix yet another test --- tests/unit/transport/drivers/test_aiokafka.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index c66deb579..146a4378a 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -1466,8 +1466,8 @@ def test__new_producer(self, *, app): "expected_args", [ pytest.param( - {"api_version": ""}, - marks=pytest.mark.conf(producer_api_version="0.10"), + {"api_version": "auto"}, + marks=pytest.mark.conf(producer_api_version="auto"), ), pytest.param({"acks": -1}, marks=pytest.mark.conf(producer_acks="all")), pytest.param( @@ -1522,7 +1522,6 @@ def test__new_producer(self, *, app): ), ], ) - @pytest.mark.skip("fix me") def test__new_producer__using_settings(self, expected_args, *, app): producer = Producer(app.transport) self.assert_new_producer(producer, **expected_args) From ea02c218058b22eacc43721957a720e90a2130f9 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 14:08:28 -0500 Subject: [PATCH 09/14] remove unneeded import --- faust/transport/drivers/aiokafka.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/faust/transport/drivers/aiokafka.py b/faust/transport/drivers/aiokafka.py index 57ff3ca0f..8e694dee4 100644 --- a/faust/transport/drivers/aiokafka.py +++ b/faust/transport/drivers/aiokafka.py @@ -26,7 +26,7 @@ import aiokafka import aiokafka.abc import opentracing -from aiokafka import AIOKafkaConsumer, TopicPartition +from aiokafka import TopicPartition from aiokafka.consumer.group_coordinator import OffsetCommitRequest from aiokafka.coordinator.assignors.roundrobin import RoundRobinPartitionAssignor from aiokafka.errors import ( @@ -830,7 +830,7 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: Returns :const:`True` if any error was logged. """ - consumer: AIOKafkaConsumer = self._ensure_consumer() + consumer = self._ensure_consumer() secs_since_started = now - self.time_started aiotp = TopicPartition(tp.topic, tp.partition) assignment = consumer._fetcher._subscriptions.subscription.assignment From 0554b005169bc271f5a795db72d7739410fbf3c1 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 15:39:16 -0500 Subject: [PATCH 10/14] re-enable another test --- tests/unit/transport/drivers/test_aiokafka.py | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 146a4378a..e12c2f703 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -1833,7 +1833,6 @@ async def test_on_start( await threaded_producer.start() await threaded_producer.stop() - @pytest.mark.skip("Needs fixing") @pytest.mark.asyncio async def test_on_thread_stop( self, From d09a8892fd3f3e863b03aa76553c24426a122d27 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 15:45:16 -0500 Subject: [PATCH 11/14] fix linting --- tests/unit/transport/drivers/test_aiokafka.py | 25 +++++++++++-------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index e12c2f703..729f8fc59 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -296,7 +296,9 @@ def _consumer(self, now, cthread, tp): ) _consumer.assignment.return_value = {tp} - _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( + ( + _consumer._fetcher._subscriptions.subscription.assignment.state_value + ).return_value = MagicMock( assignment={tp}, timestamp=now, highwater=1, @@ -539,16 +541,19 @@ def test_timed_out(self, *, cthread, now, tp, logger, _consumer): assignment = cthread.assignment() assert assignment == {tp} - _consumer._fetcher._subscriptions.subscription.assignment.state_value.return_value = MagicMock( - assignment=assignment, - timestamp=now, - highwater=None, - tp_stream_timeout_secs=cthread.tp_stream_timeout_secs, - tp_fetch_request_timeout_secs=cthread.tp_fetch_request_timeout_secs, - ) - _consumer._fetcher._subscriptions.subscription.assignment.state_value.timestamp.return_value = ( - now + fetcher = _consumer._fetcher + (fetcher._subscriptions.subscription.assignment.state_value).return_value = ( + MagicMock( + assignment=assignment, + timestamp=now, + highwater=None, + tp_stream_timeout_secs=cthread.tp_stream_timeout_secs, + tp_fetch_request_timeout_secs=cthread.tp_fetch_request_timeout_secs, + ) ) + ( + fetcher._subscriptions.subscription.assignment.state_value.timestamp + ).return_value = now assert cthread.verify_event_path(now, tp) is None logger.error.assert_called_with( From f1d6534be804d98acef190c2bf4b50e39617aeee Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 23 Feb 2024 16:02:24 -0500 Subject: [PATCH 12/14] fix another test --- tests/unit/transport/drivers/test_aiokafka.py | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/tests/unit/transport/drivers/test_aiokafka.py b/tests/unit/transport/drivers/test_aiokafka.py index 729f8fc59..d25ac8aaa 100644 --- a/tests/unit/transport/drivers/test_aiokafka.py +++ b/tests/unit/transport/drivers/test_aiokafka.py @@ -477,7 +477,6 @@ def test_timed_out(self, *, cthread, _consumer, now, tp, logger): ) -@pytest.mark.skip("Needs fixing") class Test_VEP_no_recent_fetch(Test_verify_event_path_base): def test_recent_fetch(self, *, cthread, now, tp, logger): self._set_last_response(now - 30.0) @@ -485,10 +484,15 @@ def test_recent_fetch(self, *, cthread, now, tp, logger): assert cthread.verify_event_path(now, tp) is None logger.error.assert_not_called() - def test_timed_out(self, *, cthread, now, tp, logger): + def test_timed_out(self, *, cthread, now, tp, logger, _consumer): self._set_last_response(now - 30.0) self._set_last_request(now - cthread.tp_fetch_request_timeout_secs * 2) - assert cthread.verify_event_path(now, tp) is None + assert ( + cthread.verify_event_path( + now + cthread.tp_fetch_request_timeout_secs * 2, tp + ) + is None + ) logger.error.assert_called_with( mod.SLOW_PROCESSING_NO_RECENT_FETCH, ANY, From 9c17a1726aa4bb78cbea33717f0237933165b50e Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Fri, 29 Mar 2024 20:21:37 -0400 Subject: [PATCH 13/14] Update aiokafka.py --- faust/transport/drivers/aiokafka.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/faust/transport/drivers/aiokafka.py b/faust/transport/drivers/aiokafka.py index 72e30d1e8..fe189d9d3 100644 --- a/faust/transport/drivers/aiokafka.py +++ b/faust/transport/drivers/aiokafka.py @@ -840,7 +840,7 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: poll_at = None aiotp_state = assignment.state_value(aiotp) if aiotp_state and aiotp_state.timestamp: - poll_at = aiotp_state.timestamp + poll_at = aiotp_state.timestamp / 1000 if poll_at is None: if secs_since_started >= self.tp_fetch_request_timeout_secs: # NO FETCH REQUEST SENT AT ALL SINCE WORKER START From 32e7a51316ad84aa67ed59fba5dcd7a935e79b30 Mon Sep 17 00:00:00 2001 From: William Barnhart Date: Mon, 1 Apr 2024 11:14:46 -0400 Subject: [PATCH 14/14] revert changes --- faust/transport/drivers/aiokafka.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/faust/transport/drivers/aiokafka.py b/faust/transport/drivers/aiokafka.py index 3da4d192e..72e30d1e8 100644 --- a/faust/transport/drivers/aiokafka.py +++ b/faust/transport/drivers/aiokafka.py @@ -840,14 +840,14 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: poll_at = None aiotp_state = assignment.state_value(aiotp) if aiotp_state and aiotp_state.timestamp: - poll_at = aiotp_state.timestamp / 1000 + poll_at = aiotp_state.timestamp if poll_at is None: if secs_since_started >= self.tp_fetch_request_timeout_secs: # NO FETCH REQUEST SENT AT ALL SINCE WORKER START self.log.error( SLOW_PROCESSING_NO_FETCH_SINCE_START, tp, - secs_since_started, + humanize_seconds_ago(secs_since_started), ) return True @@ -857,7 +857,7 @@ def _verify_aiokafka_event_path(self, now: float, tp: TP) -> bool: self.log.error( SLOW_PROCESSING_NO_RECENT_FETCH, tp, - secs_since_request, + humanize_seconds_ago(secs_since_request), ) return True