From 066015d59438e5525e19d46e56ae06300395781e Mon Sep 17 00:00:00 2001 From: Carter Tinney Date: Tue, 1 Sep 2026 12:36:50 -0700 Subject: [PATCH] fix: close Janus queues during async client shutdown Construct Janus 2 queues without an event-loop round trip and close every fixed and dynamic inbox after its consumers stop. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../iot/device/iothub/aio/async_clients.py | 4 ++++ .../iot/device/iothub/aio/async_inbox.py | 20 ++++++++----------- .../azure/iot/device/iothub/inbox_manager.py | 12 +++++++++++ pyproject.toml | 2 +- tests/unit/iothub/aio/test_async_clients.py | 12 +++++++++++ tests/unit/iothub/aio/test_async_inbox.py | 18 +++++++++++++++++ uv.lock | 2 +- 7 files changed, 56 insertions(+), 14 deletions(-) diff --git a/azure-iot-device/azure/iot/device/iothub/aio/async_clients.py b/azure-iot-device/azure/iot/device/iothub/aio/async_clients.py index 6eaca9c1a..959744ff7 100644 --- a/azure-iot-device/azure/iot/device/iothub/aio/async_clients.py +++ b/azure-iot-device/azure/iot/device/iothub/aio/async_clients.py @@ -6,6 +6,7 @@ """This module contains user-facing asynchronous clients for the Azure IoTHub Device SDK for Python. """ + from __future__ import annotations # Needed for annotation bug < 3.10 import logging import asyncio @@ -197,6 +198,9 @@ async def shutdown(self) -> None: # Stop the Client Event handlers now that everything else is completed self._handler_manager.stop(receiver_handlers_only=False) + # All inbox consumers have stopped, so their Janus queues can now be closed permanently. + await asyncio.gather(*(inbox.shutdown() for inbox in self._inbox_manager.get_all_inboxes())) + # Yes, that means the pipeline is disconnected twice (well, actually three times if you # consider that the client-level disconnect causes two pipeline-level disconnects for # reasons explained in comments in the client's .disconnect() method). diff --git a/azure-iot-device/azure/iot/device/iothub/aio/async_inbox.py b/azure-iot-device/azure/iot/device/iothub/aio/async_inbox.py index ec09892a0..66172c72e 100644 --- a/azure-iot-device/azure/iot/device/iothub/aio/async_inbox.py +++ b/azure-iot-device/azure/iot/device/iothub/aio/async_inbox.py @@ -4,6 +4,7 @@ # license information. # -------------------------------------------------------------------------- """This module contains an Inbox class for use with an asynchronous client""" + import asyncio import janus from azure.iot.device.iothub.sync_inbox import AbstractInbox @@ -24,18 +25,7 @@ class AsyncClientInbox(AbstractInbox): def __init__(self): """Initializer for AsyncClientInbox.""" - - # The queue must be instantiated on the client internal loop, but there's no way to do - # that at instantiation from a different loop, so instead we make coroutine to do the - # task and run it on the client internal loop. - # It's not pretty, but it works (newer versions of janus have a loop parameter, but - # not the version we are currently locked at) - async def make_queue(): - return janus.Queue() - - loop = loop_management.get_client_internal_loop() - fut = asyncio.run_coroutine_threadsafe(make_queue(), loop) - self._queue = fut.result() + self._queue = janus.Queue() def __contains__(self, item): """Return True if item is in Inbox, False otherwise""" @@ -63,6 +53,12 @@ async def get(self): fut = asyncio.run_coroutine_threadsafe(self._queue.async_q.get(), loop) return await asyncio.wrap_future(fut) + async def shutdown(self): + """Shut down the inbox and release its Janus queue resources.""" + loop = loop_management.get_client_internal_loop() + fut = asyncio.run_coroutine_threadsafe(self._queue.aclose(), loop) + await asyncio.wrap_future(fut) + def empty(self): """Returns True if the inbox is empty, False otherwise diff --git a/azure-iot-device/azure/iot/device/iothub/inbox_manager.py b/azure-iot-device/azure/iot/device/iothub/inbox_manager.py index 24056e28f..037836e2c 100644 --- a/azure-iot-device/azure/iot/device/iothub/inbox_manager.py +++ b/azure-iot-device/azure/iot/device/iothub/inbox_manager.py @@ -104,6 +104,18 @@ def get_client_event_inbox(self): """ return self.client_event_inbox + def get_all_inboxes(self): + """Retrieve every inbox managed by this instance.""" + return ( + self.unified_message_inbox, + self.generic_method_request_inbox, + self.twin_patch_inbox, + self.client_event_inbox, + self.c2d_message_inbox, + *self.input_message_inboxes.values(), + *self.named_method_request_inboxes.values(), + ) + def clear_all_method_requests(self): """Delete all method requests currently in inboxes.""" self.generic_method_request_inbox.clear() diff --git a/pyproject.toml b/pyproject.toml index ffd7676ff..37a541178 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -27,7 +27,7 @@ classifiers = [ ] dependencies = [ "deprecation>=2.1.0,<3.0.0", - "janus", + "janus>=2.0.0,<3.0.0", "paho-mqtt>=2.0.0,<3.0.0", "PySocks", "requests>=2.32.3,<3.0.0", diff --git a/tests/unit/iothub/aio/test_async_clients.py b/tests/unit/iothub/aio/test_async_clients.py index 90c4f34d5..5c3d7eb6c 100644 --- a/tests/unit/iothub/aio/test_async_clients.py +++ b/tests/unit/iothub/aio/test_async_clients.py @@ -173,6 +173,18 @@ def check_handlers_and_complete(callback): assert hm_stop_spy.call_count == 1 assert hm_stop_spy.call_args == mocker.call(receiver_handlers_only=False) + @pytest.mark.it("Shuts down all fixed and dynamically created inboxes") + async def test_shuts_down_all_inboxes(self, mocker, client): + client.disconnect = mocker.MagicMock() + client.disconnect.return_value = await create_completed_future(None) + client._inbox_manager.get_input_message_inbox("input") + client._inbox_manager.get_method_request_inbox("method") + inboxes = client._inbox_manager.get_all_inboxes() + + await client.shutdown() + + assert all(inbox._queue.closed for inbox in inboxes) + class SharedClientConnectTests(object): @pytest.mark.it("Begins a 'connect' pipeline operation") diff --git a/tests/unit/iothub/aio/test_async_inbox.py b/tests/unit/iothub/aio/test_async_inbox.py index d1da19a79..1c0236bf9 100644 --- a/tests/unit/iothub/aio/test_async_inbox.py +++ b/tests/unit/iothub/aio/test_async_inbox.py @@ -27,6 +27,17 @@ def inbox(): @pytest.mark.describe("AsyncClientInbox") class TestAsyncClientInbox(object): + @pytest.mark.it("Instantiates a Janus queue without accessing the internal event loop") + def test_instantiates_queue_without_internal_loop(self, mocker): + mock_queue_constructor = mocker.patch("azure.iot.device.iothub.aio.async_inbox.janus.Queue") + mock_get_internal_loop = mocker.patch.object(loop_management, "get_client_internal_loop") + + inbox = AsyncClientInbox() + + assert inbox._queue is mock_queue_constructor.return_value + assert mock_queue_constructor.call_args == mocker.call() + assert mock_get_internal_loop.call_count == 0 + @pytest.mark.it("Instantiates empty") def test_instantiates_empty(self, inbox): assert inbox.empty() @@ -62,6 +73,13 @@ async def test_operates_according_to_FIFO(self, mocker, inbox): assert await asyncio.wait_for(inbox.get(), timeout=PROMPT_TIMEOUT) is item2 assert await asyncio.wait_for(inbox.get(), timeout=PROMPT_TIMEOUT) is item3 + @pytest.mark.it("Closes the Janus queue on shutdown") + @pytest.mark.asyncio + async def test_shutdown_closes_queue(self, inbox): + await inbox.shutdown() + + assert inbox._queue.closed + @pytest.mark.describe("AsyncClientInbox - .put()") class TestAsyncClientInboxPut(object): diff --git a/uv.lock b/uv.lock index 4409fe282..303023b99 100644 --- a/uv.lock +++ b/uv.lock @@ -90,7 +90,7 @@ test = [ [package.metadata] requires-dist = [ { name = "deprecation", specifier = ">=2.1.0,<3.0.0" }, - { name = "janus" }, + { name = "janus", specifier = ">=2.0.0,<3.0.0" }, { name = "paho-mqtt", specifier = ">=2.0.0,<3.0.0" }, { name = "pysocks" }, { name = "requests", specifier = ">=2.32.3,<3.0.0" },