import threading import time from unittest.mock import Mock from fenrirscreenreader.core.backgroundTaskManager import BackgroundTaskManager from fenrirscreenreader.core.eventData import FenrirEventType def wait_for_call(mock, timeout=1.0): deadline = time.monotonic() + timeout while not mock.called and time.monotonic() < deadline: time.sleep(0.01) assert mock.called def test_worker_returns_result_through_event_queue_before_callback_runs(): event_manager = Mock() environment = { "runtime": { "DebugManager": Mock(), "EventManager": event_manager, } } manager = BackgroundTaskManager(worker_count=1) manager.initialize(environment) callback = Mock() main_thread_id = threading.get_ident() try: task_id = manager.submit_task(lambda value: value * 2, callback, 21) wait_for_call(event_manager.put_to_event_queue) callback.assert_not_called() event_type, result = event_manager.put_to_event_queue.call_args.args assert event_type == FenrirEventType.background_task_result assert result == { "task_id": task_id, "succeeded": True, "value": 42, "error": "", } callback.side_effect = lambda _result: setattr( callback, "thread_id", threading.get_ident() ) manager.handle_result(result) callback.assert_called_once_with(result) assert callback.thread_id == main_thread_id finally: manager.shutdown() def test_cancelled_task_result_is_not_delivered(): environment = { "runtime": { "DebugManager": Mock(), "EventManager": Mock(), } } manager = BackgroundTaskManager(worker_count=1) manager.initialize(environment) callback = Mock() try: task_id = manager.submit_task(lambda: "late", callback) manager.cancel_task(task_id) manager.handle_result( { "task_id": task_id, "succeeded": True, "value": "late", "error": "", } ) callback.assert_not_called() finally: manager.shutdown() def test_external_tasks_cannot_starve_default_tasks(): event_manager = Mock() environment = { "runtime": { "DebugManager": Mock(), "EventManager": event_manager, } } manager = BackgroundTaskManager(worker_count=1, external_worker_count=2) manager.initialize(environment) release_external = threading.Event() external_started = [threading.Event(), threading.Event()] def block_external(started): started.set() release_external.wait(timeout=2) try: for started in external_started: manager.submit_external_task(block_external, Mock(), started) for started in external_started: assert started.wait(timeout=1) manager.submit_task(lambda: "voice result", Mock()) wait_for_call(event_manager.put_to_event_queue) result = event_manager.put_to_event_queue.call_args.args[1] assert result["value"] == "voice result" finally: release_external.set() manager.shutdown()