|
@@ -152,6 +152,77 @@ class TestPrinterManager:
|
|
|
manager._schedule_async(coro)
|
|
manager._schedule_async(coro)
|
|
|
coro.close()
|
|
coro.close()
|
|
|
|
|
|
|
|
|
|
+ @staticmethod
|
|
|
|
|
+ def _run_with_loop_thread(manager, coro_fn, cancel: bool):
|
|
|
|
|
+ """Schedule on a real loop in a thread, the way the MQTT thread does.
|
|
|
|
|
+
|
|
|
|
|
+ With ``cancel``, cancel the task once it runs, as stopping the loop
|
|
|
|
|
+ does at shutdown. Returns the future ``_schedule_async`` created,
|
|
|
|
|
+ settled, with its done-callback already run.
|
|
|
|
|
+ """
|
|
|
|
|
+ import asyncio
|
|
|
|
|
+ import concurrent.futures
|
|
|
|
|
+ import threading
|
|
|
|
|
+
|
|
|
|
|
+ loop = asyncio.new_event_loop()
|
|
|
|
|
+ thread = threading.Thread(target=loop.run_forever, daemon=True)
|
|
|
|
|
+ thread.start()
|
|
|
|
|
+ futures = []
|
|
|
|
|
+ real = asyncio.run_coroutine_threadsafe
|
|
|
|
|
+
|
|
|
|
|
+ def record(coro, target_loop):
|
|
|
|
|
+ future = real(coro, target_loop)
|
|
|
|
|
+ futures.append(future)
|
|
|
|
|
+ return future
|
|
|
|
|
+
|
|
|
|
|
+ started = threading.Event()
|
|
|
|
|
+
|
|
|
|
|
+ async def wrapped():
|
|
|
|
|
+ started.set()
|
|
|
|
|
+ return await coro_fn()
|
|
|
|
|
+
|
|
|
|
|
+ try:
|
|
|
|
|
+ manager._loop = loop
|
|
|
|
|
+ with patch("asyncio.run_coroutine_threadsafe", side_effect=record):
|
|
|
|
|
+ manager._schedule_async(wrapped())
|
|
|
|
|
+ assert len(futures) == 1
|
|
|
|
|
+ if cancel:
|
|
|
|
|
+ assert started.wait(timeout=5)
|
|
|
|
|
+ loop.call_soon_threadsafe(lambda: [t.cancel() for t in asyncio.all_tasks(loop)])
|
|
|
|
|
+ concurrent.futures.wait(futures, timeout=5)
|
|
|
|
|
+ finally:
|
|
|
|
|
+ # Callbacks run on the loop thread; joining it means they have run.
|
|
|
|
|
+ loop.call_soon_threadsafe(loop.stop)
|
|
|
|
|
+ thread.join(timeout=5)
|
|
|
|
|
+ loop.close()
|
|
|
|
|
+ return futures[0]
|
|
|
|
|
+
|
|
|
|
|
+ def test_schedule_async_cancelled_callback_is_not_an_error(self, manager, caplog):
|
|
|
|
|
+ """#3243: shutdown cancels pending callbacks; that must not log an ERROR."""
|
|
|
|
|
+ import asyncio
|
|
|
|
|
+ import logging
|
|
|
|
|
+
|
|
|
|
|
+ async def slow():
|
|
|
|
|
+ await asyncio.sleep(10)
|
|
|
|
|
+
|
|
|
|
|
+ with caplog.at_level(logging.DEBUG, logger="backend.app.services.printer_manager"):
|
|
|
|
|
+ future = self._run_with_loop_thread(manager, slow, cancel=True)
|
|
|
|
|
+
|
|
|
|
|
+ assert future.cancelled()
|
|
|
|
|
+ assert "Exception in scheduled callback" not in caplog.text
|
|
|
|
|
+
|
|
|
|
|
+ def test_schedule_async_failing_callback_is_still_logged(self, manager, caplog):
|
|
|
|
|
+ import logging
|
|
|
|
|
+
|
|
|
|
|
+ async def boom():
|
|
|
|
|
+ raise RuntimeError("callback failed")
|
|
|
|
|
+
|
|
|
|
|
+ with caplog.at_level(logging.ERROR, logger="backend.app.services.printer_manager"):
|
|
|
|
|
+ future = self._run_with_loop_thread(manager, boom, cancel=False)
|
|
|
|
|
+
|
|
|
|
|
+ assert isinstance(future.exception(), RuntimeError)
|
|
|
|
|
+ assert "Exception in scheduled callback: callback failed" in caplog.text
|
|
|
|
|
+
|
|
|
def test_schedule_async_with_stopped_loop(self, manager):
|
|
def test_schedule_async_with_stopped_loop(self, manager):
|
|
|
"""Verify nothing happens when loop is not running."""
|
|
"""Verify nothing happens when loop is not running."""
|
|
|
mock_loop = MagicMock()
|
|
mock_loop = MagicMock()
|