mirror of
https://github.com/LorenEteval/Furious.git
synced 2026-10-08 22:59:48 +03:00
Fix reentrant output disposal
Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
@@ -35,6 +35,8 @@ connection policy remains outside it.
|
||||
typed exit delivery independently of log transport and rendering. The output callback runs at the GUI drain
|
||||
boundary, so bounding queue admission alone is insufficient: preserve bounded drain batches and a bounded consumer
|
||||
such as the shared log model. Test producer pressure and hidden-page draining independently.
|
||||
Consumer callbacks may synchronously dispose the output transport. End that drain turn before another queue read,
|
||||
callback or timer adjustment; the reentrant-disposal case in `tests/test_architecture_refactors.py` covers this path.
|
||||
- Parentless timers are acceptable only with a durable runtime owner and explicit disposal. Leaving the manager pool
|
||||
must not leave timers, callbacks, queues, or process handles alive. Stopping execution is not QObject destruction:
|
||||
disposal must also release monitors and output infrastructure, including for a runtime that was never started.
|
||||
|
||||
@@ -137,7 +137,7 @@ class MsgQueue(multiprocessing.queues.Queue):
|
||||
|
||||
def processMsg(self):
|
||||
"""Drain one bounded batch and adapt polling to recent activity."""
|
||||
if not callable(self.callback):
|
||||
if self._disposed or not callable(self.callback):
|
||||
return
|
||||
|
||||
hasMessages = False
|
||||
@@ -155,12 +155,18 @@ class MsgQueue(multiprocessing.queues.Queue):
|
||||
if not message.isspace():
|
||||
if messages is None:
|
||||
self.callback(message)
|
||||
|
||||
if self._disposed:
|
||||
return
|
||||
else:
|
||||
messages.append(message)
|
||||
|
||||
if messages:
|
||||
appendMany(messages)
|
||||
|
||||
if self._disposed:
|
||||
return
|
||||
|
||||
nextTimeout = (
|
||||
self.ACTIVE_DRAIN_INTERVAL
|
||||
if hasMessages
|
||||
|
||||
@@ -34,7 +34,7 @@ from Furious.Service.RuntimeLease import RuntimeEventRouter, RuntimeLease
|
||||
from PySide6 import QtCore
|
||||
from PySide6.QtNetwork import QLocalServer
|
||||
|
||||
from tests.support import isolatedSettings
|
||||
from tests.support import isolatedSettings, processQtEvents, application
|
||||
|
||||
from types import SimpleNamespace
|
||||
from unittest import TestCase, mock
|
||||
@@ -1421,6 +1421,38 @@ class ApplicationLifecycleTransactionTest(TestCase):
|
||||
messageQueue.dispose()
|
||||
self.assertIsNone(messageQueue._timerConnection)
|
||||
|
||||
def testCoreLogDrainStopsWhenItsCallbackDisposesTheQueue(self):
|
||||
"""A runtime shutdown during delivery retires the rest of this drain turn."""
|
||||
application()
|
||||
received = []
|
||||
messageQueue = None
|
||||
|
||||
def receive(message):
|
||||
received.append(message)
|
||||
messageQueue.dispose()
|
||||
|
||||
messageQueue = ProcessOutputModule.MsgQueue(msgCallback=receive)
|
||||
timerDestroyed = []
|
||||
messageQueue.timer.destroyed.connect(lambda *_args: timerDestroyed.append(True))
|
||||
|
||||
try:
|
||||
with mock.patch.object(
|
||||
messageQueue, 'getNoWait', side_effect=('first', 'second', '')
|
||||
) as get:
|
||||
messageQueue.processMsg()
|
||||
|
||||
self.assertEqual(received, ['first'])
|
||||
self.assertEqual(get.call_count, 1)
|
||||
self.assertIsNone(messageQueue.callback)
|
||||
self.assertIsNone(messageQueue._timerConnection)
|
||||
|
||||
processQtEvents()
|
||||
|
||||
self.assertEqual(timerDestroyed, [True])
|
||||
messageQueue.processMsg()
|
||||
finally:
|
||||
messageQueue.dispose()
|
||||
|
||||
def testFailuresAtMeaningfulStagesRollBackOnlyEarlierStages(self):
|
||||
expected = {
|
||||
'storage': ['cleanup storage', 'cleanup plugins'],
|
||||
|
||||
Reference in New Issue
Block a user