Fix AT-SPI queue flush deadlock

This commit is contained in:
Storm Dragon
2026-07-30 21:30:02 -04:00
parent c39f3bda47
commit ddea36c211
2 changed files with 119 additions and 28 deletions
+61 -28
View File
@@ -115,6 +115,7 @@ class EventManager:
self._churnSuppressed: bool = cthulhu_state.pauseAtspiChurn self._churnSuppressed: bool = cthulhu_state.pauseAtspiChurn
self._prioritizedContextToken: Optional[str] = cthulhu_state.prioritizedDesktopContextToken self._prioritizedContextToken: Optional[str] = cthulhu_state.prioritizedDesktopContextToken
self._desktopContextConfirmedEmpty: bool = False self._desktopContextConfirmedEmpty: bool = False
self._flushingStaleAtspiEvents: bool = False
self._relevanceBurstWindow: float = 0.15 self._relevanceBurstWindow: float = 0.15
self._relevanceBurstHistory: Dict[Tuple[str, str, str], float] = {} self._relevanceBurstHistory: Dict[Tuple[str, str, str], float] = {}
@@ -924,32 +925,32 @@ class EventManager:
return False return False
def _addToQueue(self, event: Any, asyncMode: bool) -> None: def _addToQueue(self, event: Any, asyncMode: bool) -> bool:
debugging = debug.debugEventQueue debugging = debug.debugEventQueue
if debugging: if debugging:
debug.printMessage(debug.LEVEL_ALL, " acquiring lock...") debug.printMessage(debug.LEVEL_ALL, " acquiring lock...")
self._gidleLock.acquire()
if debugging: with self._gidleLock:
debug.printMessage(debug.LEVEL_ALL, " ...acquired") effectiveAsyncMode = asyncMode or self._flushingStaleAtspiEvents
debug.printMessage(debug.LEVEL_ALL, " calling queue.put...") if debugging:
debug.printMessage(debug.LEVEL_ALL, " (full=%s)" \ debug.printMessage(debug.LEVEL_ALL, " ...acquired")
% self._eventQueue.full()) debug.printMessage(debug.LEVEL_ALL, " calling queue.put...")
debug.printMessage(debug.LEVEL_ALL, " (full=%s)" \
% self._eventQueue.full())
self._eventQueue.put(event) self._eventQueue.put(event)
if debugging: if debugging:
debug.printMessage(debug.LEVEL_ALL, " ...put complete") debug.printMessage(debug.LEVEL_ALL, " ...put complete")
if asyncMode and not self._gidleId: if effectiveAsyncMode and not self._gidleId:
if self._gilSleepTime: if self._gilSleepTime:
time.sleep(self._gilSleepTime) time.sleep(self._gilSleepTime)
self._gidleId = GLib.idle_add(self._dequeue) self._gidleId = GLib.idle_add(self._dequeue)
if debugging:
debug.printMessage(debug.LEVEL_ALL, " releasing lock...")
self._gidleLock.release()
if debug.debugEventQueue: if debug.debugEventQueue:
debug.printMessage(debug.LEVEL_ALL, " releasing lock...")
debug.printMessage(debug.LEVEL_ALL, " ...released") debug.printMessage(debug.LEVEL_ALL, " ...released")
return effectiveAsyncMode
def _queuePrintln(self, e: Any, isEnqueue: bool = True, isPrune: Optional[bool] = None) -> None: def _queuePrintln(self, e: Any, isEnqueue: bool = True, isPrune: Optional[bool] = None) -> None:
"""Convenience method to output queue-related debugging info.""" """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 = cthulhu.cthulhuApp.scriptManager.get_script(AXObject.get_application(e.source), e.source)
script.eventCache[e.type] = (e, time.time()) script.eventCache[e.type] = (e, time.time())
self._addToQueue(e, asyncMode) asyncMode = self._addToQueue(e, asyncMode)
if not asyncMode: if not asyncMode:
self._dequeue() self._dequeue()
@@ -1727,26 +1728,58 @@ class EventManager:
def _flush_stale_atspi_events(self) -> None: def _flush_stale_atspi_events(self) -> None:
"""Drops queued events that no longer match the compositor context.""" """Drops queued events that no longer match the compositor context."""
self._gidleLock.acquire() with self._gidleLock:
try: self._flushingStaleAtspiEvents = True
originalQueue = self._eventQueue originalQueue = self._eventQueue
newQueue: queue.Queue[Any] = queue.Queue(0) self._eventQueue = queue.Queue(0)
retainedEvents: list[Any] = []
try:
while not originalQueue.empty(): while not originalQueue.empty():
try: try:
event = originalQueue.get_nowait() event = originalQueue.get_nowait()
except queue.Empty: except queue.Empty:
break 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 continue
newQueue.put(event) retainedEvents.append(event)
self._eventQueue = newQueue
if self._asyncMode and not self._eventQueue.empty() and not self._gidleId:
self._gidleId = GLib.idle_add(self._dequeue)
finally: 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: def _inFlood(self) -> bool:
size = self._eventQueue.qsize() size = self._eventQueue.qsize()
@@ -186,6 +186,64 @@ class EventManagerCompositorContextRegressionTests(unittest.TestCase):
self.assertEqual(list(self.manager._eventQueue.queue), [currentEvent]) 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: def test_stale_background_event_does_not_activate_script_during_suppression(self) -> None:
script = mock.Mock() script = mock.Mock()
script.isActivatableEvent.return_value = True script.isActivatableEvent.return_value = True