From b802fa57015c0d3383d2b67ffe6f8499fae23d76 Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 15:52:21 -0400 Subject: [PATCH 1/6] fix: remove expired ack_ids --- .../pubsub_v1/subscriber/_protocol/leaser.py | 7 +- .../_protocol/streaming_pull_manager.py | 47 +++++---- .../subscriber/test_streaming_pull_manager.py | 97 +++++++++++++++++-- 3 files changed, 124 insertions(+), 27 deletions(-) diff --git a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py index 508f4d7ce..968b15515 100644 --- a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py +++ b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py @@ -197,7 +197,12 @@ def maintain_leases(self) -> None: # is inactive. assert self._manager.dispatcher is not None ack_id_gen = (ack_id for ack_id in ack_ids) - self._manager._send_lease_modacks(ack_id_gen, deadline) + expired_ack_ids = self._manager._send_lease_modacks( + ack_id_gen, deadline + ) + if self._manager._exactly_once_delivery_enabled(): + for ack_id in expired_ack_ids: + leased_messages.pop(ack_id) # Now wait an appropriate period of time and do this again. # diff --git a/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py b/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py index 932699261..21c1bab7b 100644 --- a/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py +++ b/google/cloud/pubsub_v1/subscriber/_protocol/streaming_pull_manager.py @@ -988,7 +988,9 @@ def _get_initial_request( # Return the initial request. return request - def _send_lease_modacks(self, ack_ids: Iterable[str], ack_deadline: float): + def _send_lease_modacks( + self, ack_ids: Iterable[str], ack_deadline: float + ) -> List[str]: exactly_once_enabled = False with self._exactly_once_enabled_lock: exactly_once_enabled = self._exactly_once_enabled @@ -1002,15 +1004,19 @@ def _send_lease_modacks(self, ack_ids: Iterable[str], ack_deadline: float): assert self._dispatcher is not None self._dispatcher.modify_ack_deadline(items) + expired_ack_ids = [] for req in items: try: assert req.future is not None req.future.result() - except AcknowledgeError: + except AcknowledgeError as ack_error: _LOGGER.warning( "AcknowledgeError when lease-modacking a message.", exc_info=True, ) + if ack_error.error_code == AcknowledgeStatus.INVALID_ACK_ID: + expired_ack_ids.append(req.ack_id) + return expired_ack_ids else: items = [ requests.ModAckRequest(ack_id, self.ack_deadline, None) @@ -1018,6 +1024,7 @@ def _send_lease_modacks(self, ack_ids: Iterable[str], ack_deadline: float): ] assert self._dispatcher is not None self._dispatcher.modify_ack_deadline(items) + return [] def _exactly_once_delivery_enabled(self) -> bool: """Whether exactly-once delivery is enabled for the subscription.""" @@ -1071,28 +1078,32 @@ def _on_response(self, response: gapic_types.StreamingPullResponse) -> None: # modack the messages we received, as this tells the server that we've # received them. ack_id_gen = (message.ack_id for message in received_messages) - self._send_lease_modacks(ack_id_gen, self.ack_deadline) + expired_ack_ids = set(self._send_lease_modacks(ack_id_gen, self.ack_deadline)) with self._pause_resume_lock: assert self._scheduler is not None assert self._leaser is not None for received_message in received_messages: - message = google.cloud.pubsub_v1.subscriber.message.Message( - received_message.message, - received_message.ack_id, - received_message.delivery_attempt, - self._scheduler.queue, - self._exactly_once_delivery_enabled, - ) - self._messages_on_hold.put(message) - self._on_hold_bytes += message.size - req = requests.LeaseRequest( - ack_id=message.ack_id, - byte_size=message.size, - ordering_key=message.ordering_key, - ) - self._leaser.add([req]) + if ( + not self._exactly_once_delivery_enabled() + or received_message.ack_id not in expired_ack_ids + ): + message = google.cloud.pubsub_v1.subscriber.message.Message( + received_message.message, + received_message.ack_id, + received_message.delivery_attempt, + self._scheduler.queue, + self._exactly_once_delivery_enabled, + ) + self._messages_on_hold.put(message) + self._on_hold_bytes += message.size + req = requests.LeaseRequest( + ack_id=message.ack_id, + byte_size=message.size, + ordering_key=message.ordering_key, + ) + self._leaser.add([req]) self._maybe_release_messages() diff --git a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py index deb476eb1..c6c35aca1 100644 --- a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py +++ b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py @@ -1069,6 +1069,60 @@ def test_send_unary_modack_retry_error_exactly_once_enabled_with_futures( ) +@mock.patch("google.api_core.bidi.ResumableBidiRpc", autospec=True) +@mock.patch("google.api_core.bidi.BackgroundConsumer", autospec=True) +@mock.patch("google.cloud.pubsub_v1.subscriber._protocol.leaser.Leaser", autospec=True) +@mock.patch( + "google.cloud.pubsub_v1.subscriber._protocol.dispatcher.Dispatcher", autospec=True +) +@mock.patch( + "google.cloud.pubsub_v1.subscriber._protocol.heartbeater.Heartbeater", autospec=True +) +def test_open(heartbeater, dispatcher, leaser, background_consumer, resumable_bidi_rpc): + manager = make_manager() + + with mock.patch.object( + type(manager), "ack_deadline", new=mock.PropertyMock(return_value=18) + ): + manager.open(mock.sentinel.callback, mock.sentinel.on_callback_error) + + heartbeater.assert_called_once_with(manager) + heartbeater.return_value.start.assert_called_once() + assert manager._heartbeater == heartbeater.return_value + + dispatcher.assert_called_once_with(manager, manager._scheduler.queue) + dispatcher.return_value.start.assert_called_once() + assert manager._dispatcher == dispatcher.return_value + + leaser.assert_called_once_with(manager) + leaser.return_value.start.assert_called_once() + assert manager.leaser == leaser.return_value + + background_consumer.assert_called_once_with(manager._rpc, manager._on_response) + background_consumer.return_value.start.assert_called_once() + assert manager._consumer == background_consumer.return_value + + resumable_bidi_rpc.assert_called_once_with( + start_rpc=manager._client.streaming_pull, + initial_request=mock.ANY, + should_recover=manager._should_recover, + should_terminate=manager._should_terminate, + throttle_reopen=True, + ) + initial_request_arg = resumable_bidi_rpc.call_args.kwargs["initial_request"] + assert initial_request_arg.func == manager._get_initial_request + assert initial_request_arg.args[0] == 60 + assert not manager._client.get_subscription.called + + resumable_bidi_rpc.return_value.add_done_callback.assert_called_once_with( + manager._on_rpc_done + ) + assert manager._rpc == resumable_bidi_rpc.return_value + + manager._consumer.is_active = True + assert manager.is_active is True + + def test_heartbeat(): manager = make_manager() manager._rpc = mock.create_autospec(bidi.BidiRpc, instance=True) @@ -1120,6 +1174,7 @@ def test_heartbeat_stream_ack_deadline_seconds(caplog): "google.cloud.pubsub_v1.subscriber._protocol.heartbeater.Heartbeater", autospec=True ) def test_open(heartbeater, dispatcher, leaser, background_consumer, resumable_bidi_rpc): + manager = make_manager() with mock.patch.object( @@ -1847,11 +1902,18 @@ def test__on_response_exactly_once_immediate_modacks_fail(): def complete_futures_with_error(*args, **kwargs): modack_requests = args[0] for req in modack_requests: - req.future.set_exception( - subscriber_exceptions.AcknowledgeError( - subscriber_exceptions.AcknowledgeStatus.SUCCESS, None + if req.ack_id == "fack": + req.future.set_exception( + subscriber_exceptions.AcknowledgeError( + subscriber_exceptions.AcknowledgeStatus.INVALID_ACK_ID, None + ) + ) + else: + req.future.set_exception( + subscriber_exceptions.AcknowledgeError( + subscriber_exceptions.AcknowledgeStatus.SUCCESS, None + ) ) - ) dispatcher.modify_ack_deadline.side_effect = complete_futures_with_error @@ -1861,19 +1923,38 @@ def complete_futures_with_error(*args, **kwargs): gapic_types.ReceivedMessage( ack_id="fack", message=gapic_types.PubsubMessage(data=b"foo", message_id="1"), - ) + ), + gapic_types.ReceivedMessage( + ack_id="good", + message=gapic_types.PubsubMessage(data=b"foo", message_id="2"), + ), ], subscription_properties=gapic_types.StreamingPullResponse.SubscriptionProperties( exactly_once_delivery_enabled=True ), ) - # adjust message bookkeeping in leaser - fake_leaser_add(leaser, init_msg_count=0, assumed_msg_size=42) + # Actually run the method and prove that modack and schedule are called in + # the expected way. + manager._on_response(response) + + # The second messages should be scheduled, and not the first. + + schedule_calls = scheduler.schedule.mock_calls + assert len(schedule_calls) == 1 + call_args = schedule_calls[0][1] + assert call_args[0] == mock.sentinel.callback + assert isinstance(call_args[1], message.Message) + assert call_args[1].message_id == "2" + + assert manager._messages_on_hold.size == 0 + # No messages available + assert manager._messages_on_hold.get() is None # exactly_once should be enabled manager._on_response(response) - # exceptions are logged, but otherwise no effect + # do not add message + assert manager.load == 0 def test__should_recover_true(): From f33e0436320563c9b8e53947d0d00ebc3934f7df Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 16:00:36 -0400 Subject: [PATCH 2/6] remove extra test --- .../subscriber/test_streaming_pull_manager.py | 54 ------------------- 1 file changed, 54 deletions(-) diff --git a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py index c6c35aca1..0d6273c80 100644 --- a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py +++ b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py @@ -1069,60 +1069,6 @@ def test_send_unary_modack_retry_error_exactly_once_enabled_with_futures( ) -@mock.patch("google.api_core.bidi.ResumableBidiRpc", autospec=True) -@mock.patch("google.api_core.bidi.BackgroundConsumer", autospec=True) -@mock.patch("google.cloud.pubsub_v1.subscriber._protocol.leaser.Leaser", autospec=True) -@mock.patch( - "google.cloud.pubsub_v1.subscriber._protocol.dispatcher.Dispatcher", autospec=True -) -@mock.patch( - "google.cloud.pubsub_v1.subscriber._protocol.heartbeater.Heartbeater", autospec=True -) -def test_open(heartbeater, dispatcher, leaser, background_consumer, resumable_bidi_rpc): - manager = make_manager() - - with mock.patch.object( - type(manager), "ack_deadline", new=mock.PropertyMock(return_value=18) - ): - manager.open(mock.sentinel.callback, mock.sentinel.on_callback_error) - - heartbeater.assert_called_once_with(manager) - heartbeater.return_value.start.assert_called_once() - assert manager._heartbeater == heartbeater.return_value - - dispatcher.assert_called_once_with(manager, manager._scheduler.queue) - dispatcher.return_value.start.assert_called_once() - assert manager._dispatcher == dispatcher.return_value - - leaser.assert_called_once_with(manager) - leaser.return_value.start.assert_called_once() - assert manager.leaser == leaser.return_value - - background_consumer.assert_called_once_with(manager._rpc, manager._on_response) - background_consumer.return_value.start.assert_called_once() - assert manager._consumer == background_consumer.return_value - - resumable_bidi_rpc.assert_called_once_with( - start_rpc=manager._client.streaming_pull, - initial_request=mock.ANY, - should_recover=manager._should_recover, - should_terminate=manager._should_terminate, - throttle_reopen=True, - ) - initial_request_arg = resumable_bidi_rpc.call_args.kwargs["initial_request"] - assert initial_request_arg.func == manager._get_initial_request - assert initial_request_arg.args[0] == 60 - assert not manager._client.get_subscription.called - - resumable_bidi_rpc.return_value.add_done_callback.assert_called_once_with( - manager._on_rpc_done - ) - assert manager._rpc == resumable_bidi_rpc.return_value - - manager._consumer.is_active = True - assert manager.is_active is True - - def test_heartbeat(): manager = make_manager() manager._rpc = mock.create_autospec(bidi.BidiRpc, instance=True) From 1023618a53548e6526a3777e9696de449a23334a Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 16:25:50 -0400 Subject: [PATCH 3/6] fix spm test --- .../pubsub_v1/subscriber/test_streaming_pull_manager.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py index 0d6273c80..b4d841177 100644 --- a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py +++ b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py @@ -1882,6 +1882,9 @@ def complete_futures_with_error(*args, **kwargs): # Actually run the method and prove that modack and schedule are called in # the expected way. + + fake_leaser_add(leaser, init_msg_count=0, assumed_msg_size=10) + manager._on_response(response) # The second messages should be scheduled, and not the first. @@ -1897,10 +1900,8 @@ def complete_futures_with_error(*args, **kwargs): # No messages available assert manager._messages_on_hold.get() is None - # exactly_once should be enabled - manager._on_response(response) # do not add message - assert manager.load == 0 + assert manager.load == .001 def test__should_recover_true(): From 152a23f1fc8cd8939335df73361d8353621c3270 Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 16:27:51 -0400 Subject: [PATCH 4/6] fix test lint --- .../pubsub_v1/subscriber/test_streaming_pull_manager.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py index b4d841177..02d042904 100644 --- a/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py +++ b/tests/unit/pubsub_v1/subscriber/test_streaming_pull_manager.py @@ -1882,9 +1882,9 @@ def complete_futures_with_error(*args, **kwargs): # Actually run the method and prove that modack and schedule are called in # the expected way. - + fake_leaser_add(leaser, init_msg_count=0, assumed_msg_size=10) - + manager._on_response(response) # The second messages should be scheduled, and not the first. @@ -1901,7 +1901,7 @@ def complete_futures_with_error(*args, **kwargs): assert manager._messages_on_hold.get() is None # do not add message - assert manager.load == .001 + assert manager.load == 0.001 def test__should_recover_true(): From 2396055dd2f76ace79ec796c20172edcd7469ab1 Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 16:39:57 -0400 Subject: [PATCH 5/6] adding extra dropping logic --- .../pubsub_v1/subscriber/_protocol/leaser.py | 25 ++++++++++++++++--- 1 file changed, 22 insertions(+), 3 deletions(-) diff --git a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py index 968b15515..f8ef5b85e 100644 --- a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py +++ b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py @@ -186,6 +186,7 @@ def maintain_leases(self) -> None: # Create a modack request. # We do not actually call `modify_ack_deadline` over and over # because it is more efficient to make a single request. + expired_ack_ids = [] ack_ids = leased_messages.keys() if ack_ids: _LOGGER.debug("Renewing lease for %d ack IDs.", len(ack_ids)) @@ -200,9 +201,27 @@ def maintain_leases(self) -> None: expired_ack_ids = self._manager._send_lease_modacks( ack_id_gen, deadline ) - if self._manager._exactly_once_delivery_enabled(): - for ack_id in expired_ack_ids: - leased_messages.pop(ack_id) + + # Drop more items based on expiration: + if self._manager._exactly_once_delivery_enabled(): + to_drop = [ + requests.DropRequest(ack_id, item.size, item.ordering_key) + for ack_id, item in leased_messages.items() + if ack_id in expired_ack_ids + ] + if to_drop: + _LOGGER.warning( + "Dropping %s items because they were leased too long.", + len(to_drop), + ) + assert self._manager.dispatcher is not None + self._manager.dispatcher.drop(to_drop) + + # Remove dropped items from our copy of the leased messages (they + # have already been removed from the real one by + # self._manager.drop(), which calls self.remove()). + for item in to_drop: + leased_messages.pop(item.ack_id) # Now wait an appropriate period of time and do this again. # From d6055a6c0ca409f9854742fa5fc0d99f4b4f672e Mon Sep 17 00:00:00 2001 From: acocuzzo Date: Wed, 21 Sep 2022 16:40:39 -0400 Subject: [PATCH 6/6] reverting leaser.py --- .../pubsub_v1/subscriber/_protocol/leaser.py | 26 +------------------ 1 file changed, 1 insertion(+), 25 deletions(-) diff --git a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py index f8ef5b85e..508f4d7ce 100644 --- a/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py +++ b/google/cloud/pubsub_v1/subscriber/_protocol/leaser.py @@ -186,7 +186,6 @@ def maintain_leases(self) -> None: # Create a modack request. # We do not actually call `modify_ack_deadline` over and over # because it is more efficient to make a single request. - expired_ack_ids = [] ack_ids = leased_messages.keys() if ack_ids: _LOGGER.debug("Renewing lease for %d ack IDs.", len(ack_ids)) @@ -198,30 +197,7 @@ def maintain_leases(self) -> None: # is inactive. assert self._manager.dispatcher is not None ack_id_gen = (ack_id for ack_id in ack_ids) - expired_ack_ids = self._manager._send_lease_modacks( - ack_id_gen, deadline - ) - - # Drop more items based on expiration: - if self._manager._exactly_once_delivery_enabled(): - to_drop = [ - requests.DropRequest(ack_id, item.size, item.ordering_key) - for ack_id, item in leased_messages.items() - if ack_id in expired_ack_ids - ] - if to_drop: - _LOGGER.warning( - "Dropping %s items because they were leased too long.", - len(to_drop), - ) - assert self._manager.dispatcher is not None - self._manager.dispatcher.drop(to_drop) - - # Remove dropped items from our copy of the leased messages (they - # have already been removed from the real one by - # self._manager.drop(), which calls self.remove()). - for item in to_drop: - leased_messages.pop(item.ack_id) + self._manager._send_lease_modacks(ack_id_gen, deadline) # Now wait an appropriate period of time and do this again. #