diff --git a/src/cthulhu/event_manager.py b/src/cthulhu/event_manager.py index 332696f..55d9eff 100644 --- a/src/cthulhu/event_manager.py +++ b/src/cthulhu/event_manager.py @@ -115,6 +115,7 @@ class EventManager: self._churnSuppressed: bool = cthulhu_state.pauseAtspiChurn self._prioritizedContextToken: Optional[str] = cthulhu_state.prioritizedDesktopContextToken self._desktopContextConfirmedEmpty: bool = False + self._flushingStaleAtspiEvents: bool = False self._relevanceBurstWindow: float = 0.15 self._relevanceBurstHistory: Dict[Tuple[str, str, str], float] = {} @@ -924,32 +925,32 @@ class EventManager: return False - def _addToQueue(self, event: Any, asyncMode: bool) -> None: + def _addToQueue(self, event: Any, asyncMode: bool) -> bool: debugging = debug.debugEventQueue if debugging: debug.printMessage(debug.LEVEL_ALL, " acquiring lock...") - self._gidleLock.acquire() - if debugging: - debug.printMessage(debug.LEVEL_ALL, " ...acquired") - debug.printMessage(debug.LEVEL_ALL, " calling queue.put...") - debug.printMessage(debug.LEVEL_ALL, " (full=%s)" \ - % self._eventQueue.full()) + with self._gidleLock: + effectiveAsyncMode = asyncMode or self._flushingStaleAtspiEvents + if debugging: + debug.printMessage(debug.LEVEL_ALL, " ...acquired") + debug.printMessage(debug.LEVEL_ALL, " calling queue.put...") + debug.printMessage(debug.LEVEL_ALL, " (full=%s)" \ + % self._eventQueue.full()) - self._eventQueue.put(event) - if debugging: - debug.printMessage(debug.LEVEL_ALL, " ...put complete") + self._eventQueue.put(event) + if debugging: + debug.printMessage(debug.LEVEL_ALL, " ...put complete") - if asyncMode and not self._gidleId: - if self._gilSleepTime: - time.sleep(self._gilSleepTime) - self._gidleId = GLib.idle_add(self._dequeue) + if effectiveAsyncMode and not self._gidleId: + if self._gilSleepTime: + time.sleep(self._gilSleepTime) + self._gidleId = GLib.idle_add(self._dequeue) - if debugging: - debug.printMessage(debug.LEVEL_ALL, " releasing lock...") - self._gidleLock.release() if debug.debugEventQueue: + debug.printMessage(debug.LEVEL_ALL, " releasing lock...") debug.printMessage(debug.LEVEL_ALL, " ...released") + return effectiveAsyncMode def _queuePrintln(self, e: Any, isEnqueue: bool = True, isPrune: Optional[bool] = None) -> None: """Convenience method to output queue-related debugging info.""" @@ -1118,7 +1119,7 @@ class EventManager: script = cthulhu.cthulhuApp.scriptManager.get_script(AXObject.get_application(e.source), e.source) script.eventCache[e.type] = (e, time.time()) - self._addToQueue(e, asyncMode) + asyncMode = self._addToQueue(e, asyncMode) if not asyncMode: self._dequeue() @@ -1727,26 +1728,58 @@ class EventManager: def _flush_stale_atspi_events(self) -> None: """Drops queued events that no longer match the compositor context.""" - self._gidleLock.acquire() - try: + with self._gidleLock: + self._flushingStaleAtspiEvents = True originalQueue = self._eventQueue - newQueue: queue.Queue[Any] = queue.Queue(0) + self._eventQueue = queue.Queue(0) + + retainedEvents: list[Any] = [] + try: while not originalQueue.empty(): try: event = originalQueue.get_nowait() except queue.Empty: break - if self._event_is_from_stale_context(event) and not self._should_preserve_during_suppression(event): + # Accessible property calls can dispatch nested AT-SPI events. + # Keep them outside _gidleLock so a nested enqueue cannot deadlock. + try: + isStale = self._event_is_from_stale_context(event) + shouldPreserve = ( + isStale + and self._should_preserve_during_suppression(event) + ) + except Exception: + retainedEvents.append(event) + raise + + if isStale and not shouldPreserve: continue - newQueue.put(event) - - self._eventQueue = newQueue - if self._asyncMode and not self._eventQueue.empty() and not self._gidleId: - self._gidleId = GLib.idle_add(self._dequeue) + retainedEvents.append(event) finally: - self._gidleLock.release() + while not originalQueue.empty(): + try: + retainedEvents.append(originalQueue.get_nowait()) + except queue.Empty: + break + + with self._gidleLock: + eventsQueuedDuringFlush = self._eventQueue + mergedQueue: queue.Queue[Any] = queue.Queue(0) + for event in retainedEvents: + mergedQueue.put(event) + + while not eventsQueuedDuringFlush.empty(): + try: + mergedQueue.put(eventsQueuedDuringFlush.get_nowait()) + except queue.Empty: + break + + self._eventQueue = mergedQueue + self._flushingStaleAtspiEvents = False + if self._asyncMode and not self._eventQueue.empty() and not self._gidleId: + self._gidleId = GLib.idle_add(self._dequeue) def _inFlood(self) -> bool: size = self._eventQueue.qsize() diff --git a/tests/test_event_manager_compositor_context_regressions.py b/tests/test_event_manager_compositor_context_regressions.py index 587f69b..b6d0b0d 100644 --- a/tests/test_event_manager_compositor_context_regressions.py +++ b/tests/test_event_manager_compositor_context_regressions.py @@ -186,6 +186,64 @@ class EventManagerCompositorContextRegressionTests(unittest.TestCase): self.assertEqual(list(self.manager._eventQueue.queue), [currentEvent]) + def test_flush_classifies_outside_queue_lock_and_preserves_nested_event(self) -> None: + staleEvent = FakeEvent("object:children-changed:add", source="stale") + currentEvent = FakeEvent("object:children-changed:add", source="current") + nestedEvent = FakeEvent("object:state-changed:focused", source="nested", detail1=1) + self.manager._eventQueue.put(staleEvent) + self.manager._eventQueue.put(currentEvent) + self.manager._ignore = mock.Mock(return_value=False) + self.manager._prioritizeSelfHostedFocusedEvent = mock.Mock(return_value=False) + self.manager._queuePrintln = mock.Mock() + self.manager._inFlood = mock.Mock(return_value=False) + self.manager._shouldSuspendEventsFor = mock.Mock(return_value=False) + self.manager._dequeue = mock.Mock(return_value=False) + script = types.SimpleNamespace(eventCache={}) + + def classify_event(event) -> bool: + self.assertFalse(self.manager._gidleLock.locked()) + if event is staleEvent: + self.manager._enqueue(nestedEvent) + return True + return False + + self.manager._event_is_from_stale_context = mock.Mock(side_effect=classify_event) + self.manager._should_preserve_during_suppression = mock.Mock(return_value=False) + + with ( + mock.patch.object(event_manager.AXObject, "get_application", return_value=object()), + mock.patch.object(event_manager.GLib, "idle_add", return_value=1), + mock.patch.object( + event_manager.cthulhu.cthulhuApp.scriptManager, + "get_script", + return_value=script, + ), + mock.patch.object( + event_manager.AXUtilities, + "get_application_toolkit_name", + return_value="VCL", + ), + ): + self.manager._flush_stale_atspi_events() + + self.assertEqual(list(self.manager._eventQueue.queue), [currentEvent, nestedEvent]) + self.manager._dequeue.assert_not_called() + + def test_flush_restores_retained_and_unexamined_events_after_exception(self) -> None: + firstEvent = FakeEvent("object:children-changed:add", source="first") + secondEvent = FakeEvent("object:children-changed:add", source="second") + self.manager._eventQueue.put(firstEvent) + self.manager._eventQueue.put(secondEvent) + self.manager._event_is_from_stale_context = mock.Mock( + side_effect=RuntimeError("property lookup failed"), + ) + + with self.assertRaisesRegex(RuntimeError, "property lookup failed"): + self.manager._flush_stale_atspi_events() + + self.assertEqual(list(self.manager._eventQueue.queue), [firstEvent, secondEvent]) + self.assertFalse(self.manager._flushingStaleAtspiEvents) + def test_stale_background_event_does_not_activate_script_during_suppression(self) -> None: script = mock.Mock() script.isActivatableEvent.return_value = True