Bound and adapt core output draining

Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
Loren Eteval
2026-08-23 11:26:20 +08:00
parent 6750cd0402
commit 45ddb3a725
5 changed files with 135 additions and 64 deletions
+2 -4
View File
@@ -100,10 +100,8 @@ def startXrayCore(jsonString: str, msgQueue: multiprocessing.Queue):
for line in iter(file.readline, b''):
if line and not line.isspace():
try:
msgQueue.put_nowait(line.decode('utf-8', 'replace'))
except Exception:
# Any non-exit exceptions
msgQueue.putMessage(line.decode('utf-8', 'replace'))
except (OSError, ValueError):
pass
time.sleep(MsgQueue.MSG_PRODUCE_THRESHOLD / 1000)
+4
View File
@@ -6,6 +6,10 @@
terminate/join/kill escalation, reaps the child, stops timers/queues, clears callbacks, and is idempotent.
- Preserve platform exit codes and expose actionable startup errors. Do not convert serialization/start failures to an
unexplained success or “Unknown error” when context exists.
- Child-output transports are bounded and non-blocking for producers. Drain them in bounded batches regardless of page
visibility: use a short interval while messages flow and back off to a finite maximum interval while idle. Truncate
oversized messages and drop excess burst output at the documented queue boundary rather than retaining it
indefinitely. Presentation may remain lazy after collection.
- Process targets do not touch GUI objects. Queue/monitor callbacks cross to the owning Qt thread through the
established timer/signal boundary.
- A timer without a QObject parent requires an explicit durable Python owner and `dispose()` path. Never rely on wrapper
+65 -57
View File
@@ -33,6 +33,7 @@ import os
import sys
import time
import uuid
import queue
import logging
import threading
import multiprocessing
@@ -117,24 +118,34 @@ class CoreLaunchSpec:
class MsgQueue(multiprocessing.queues.Queue):
"""Deliver child-process log messages to Qt callbacks at an adaptive rate."""
"""Continuously drain a bounded child-process log queue into a callback."""
MSG_PRODUCE_THRESHOLD = 1024
OPTIMIZER_MIN_FREQ = 2
OPTIMIZER_MAX_FREQ = 256
ACTIVE_DRAIN_INTERVAL = 16
MAXIMUM_IDLE_DRAIN_INTERVAL = 256
DRAIN_INTERVAL = ACTIVE_DRAIN_INTERVAL
MAXIMUM_PENDING_MESSAGES = 1024
MAXIMUM_MESSAGE_CHARACTERS = 64 * 1024
MAXIMUM_MESSAGES_PER_TICK = 512
TRUNCATION_MARKER = '\n... [core log message truncated]'
def __init__(self, **kwargs):
"""Initialize the MsgQueue."""
msgCallback = kwargs.pop('msgCallback', None)
backgroundOptimizer = kwargs.pop('backgroundOptimizer', None)
msgCallback, maximumPendingMessages = (
kwargs.pop('msgCallback', None),
kwargs.pop('maximumPendingMessages', self.MAXIMUM_PENDING_MESSAGES),
)
super().__init__(**kwargs, ctx=multiprocessing.get_context())
super().__init__(
maxsize=maximumPendingMessages,
**kwargs,
ctx=multiprocessing.get_context(),
)
self.timer = QtCore.QTimer()
self.timer.timeout.connect(self.processMsg)
self.timeout = MsgQueue.MSG_PRODUCE_THRESHOLD
self.timeout = self.ACTIVE_DRAIN_INTERVAL
self.callback = msgCallback
self.backgroundOptimizer = backgroundOptimizer
def getNoWait(self) -> str:
"""Return no wait."""
@@ -172,57 +183,60 @@ class MsgQueue(multiprocessing.queues.Queue):
self.timer.deleteLater()
self.callback = None
self.backgroundOptimizer = None
try:
self.close()
except (OSError, ValueError):
pass
@property
def optimizer(self):
"""Return the optimizer value."""
try:
return self.backgroundOptimizer()
except Exception:
# Any non-exit exceptions
def putMessage(self, message) -> bool:
"""Queue one bounded message without ever blocking a core process."""
text = str(message)
return None
if len(text) > self.MAXIMUM_MESSAGE_CHARACTERS:
mark = self.TRUNCATION_MARKER
text = text[: self.MAXIMUM_MESSAGE_CHARACTERS - len(mark)] + mark
try:
self.put_nowait(text)
except queue.Full:
return False
except (OSError, ValueError):
return False
return True
def processMsg(self):
"""Process msg."""
msg = self.getNoWait()
"""Drain one bounded batch and adapt polling to recent queue activity."""
if not callable(self.callback):
# Nothing to do
return
if msg and not msg.isspace():
# Call message callback
self.callback(msg)
hasMessages = False
if self.optimizer is not None and self.optimizer.isVisible():
# Log page is visible: maximum loading speed for user
for _ in range(self.MAXIMUM_MESSAGES_PER_TICK):
msg = self.getNoWait()
# 1024, 512, 256, 128, 64, 32, 16, 8, 4, 2, 2, 2, ...
# For timeout value 2 Furious can handle at about 500 messages per second
self.setTimeout(
max(MsgQueue.OPTIMIZER_MIN_FREQ, self.getTimeout() // 2)
)
self.startTimer()
else:
# Log page is unavailable or hidden: low speed in background
# to avoid consuming too much CPU resources
if not msg:
break
# 1024, 512, 256, 256, 256, ...
# For timeout value 256 Furious can handle at about 4 messages per second
self.setTimeout(
max(MsgQueue.OPTIMIZER_MAX_FREQ, self.getTimeout() // 2)
)
self.startTimer()
hasMessages = True
if not msg.isspace():
self.callback(msg)
if hasMessages:
nextTimeout = self.ACTIVE_DRAIN_INTERVAL
else:
# Reset timeout value
self.setTimeout(MsgQueue.MSG_PRODUCE_THRESHOLD)
nextTimeout = min(
self.MAXIMUM_IDLE_DRAIN_INTERVAL,
max(
self.ACTIVE_DRAIN_INTERVAL,
self.getTimeout() * 2,
),
)
if nextTimeout != self.getTimeout():
self.setTimeout(nextTimeout)
self.startTimer()
@@ -306,7 +320,9 @@ class CoreProcessMonitor(CoreRuntime, ABC):
try:
self.process.close()
except Exception:
# close() can fail if the process handle is still considered active.
# Any non-exit exceptions
# close() can fail while the handle is active.
pass
self.process = None
@@ -331,14 +347,10 @@ class CoreProcessWorker(CoreProcessMonitor, ABC):
def __init__(self, **kwargs):
"""Initialize the CoreProcessWorker."""
msgCallback = kwargs.pop('msgCallback', None)
# Drain output more frequently while the unified log page is visible.
backgroundOptimizer = kwargs.pop('backgroundOptimizer', AppLogPage)
super().__init__(**kwargs)
self.msgQueue = MsgQueue(
msgCallback=msgCallback, backgroundOptimizer=backgroundOptimizer
)
self.msgQueue = MsgQueue(msgCallback=msgCallback)
def handleInternalProcessStopped(self):
"""Handle internal process stopped."""
@@ -408,7 +420,7 @@ class CoreProcessWorker(CoreProcessMonitor, ABC):
logger.info(f'{self.name()} {self.version()} started')
self.msgQueue.setTimeout(MsgQueue.MSG_PRODUCE_THRESHOLD)
self.msgQueue.setTimeout(MsgQueue.ACTIVE_DRAIN_INTERVAL)
self.msgQueue.startTimer()
if launchSpec.waitCore:
@@ -463,9 +475,7 @@ class ProcessOutputRedirector:
TemporaryDir = QtCore.QTemporaryDir()
@staticmethod
def launch(
msgQueue: multiprocessing.Queue, entrypoint: Callable[[], None], redirect: bool
):
def launch(msgQueue: MsgQueue, entrypoint: Callable[[], None], redirect: bool):
"""Run an entry point while forwarding its output to a message queue."""
if not callable(entrypoint):
return
@@ -503,10 +513,8 @@ class ProcessOutputRedirector:
for line in iter(file.readline, b''):
if line and not line.isspace():
try:
msgQueue.put_nowait(line.decode('utf-8', 'replace'))
except Exception:
# Any non-exit exceptions
msgQueue.putMessage(line.decode('utf-8', 'replace'))
except (OSError, ValueError):
pass
time.sleep(MsgQueue.MSG_PRODUCE_THRESHOLD / 1000)
+1 -3
View File
@@ -60,9 +60,7 @@ class Tun2socks(CoreProcessWorker):
def __init__(self, **kwargs):
"""Initialize the Tun2socks."""
backgroundOptimizer = kwargs.pop('backgroundOptimizer', AppLogPage)
super().__init__(**kwargs, backgroundOptimizer=backgroundOptimizer)
super().__init__(**kwargs)
self.cleanup = None
+63
View File
@@ -361,6 +361,69 @@ class ApplicationLifecycleTransactionTest(TestCase):
thread.start.assert_called_once_with()
entrypoint.assert_called_once_with()
def testCoreLogQueueIsBoundedAndTruncatesBeforeTransport(self):
"""Drop excess burst output instead of retaining an unbounded backlog."""
received = []
messageQueue = CoreProcessWorkerModule.MsgQueue(
msgCallback=received.append,
maximumPendingMessages=4,
)
try:
oversized = 'x' * (
CoreProcessWorkerModule.MsgQueue.MAXIMUM_MESSAGE_CHARACTERS + 100
)
accepted = [
messageQueue.putMessage(oversized if index == 0 else str(index))
for index in range(20)
]
self.assertLessEqual(sum(accepted), 4)
for _ in range(20):
messageQueue.processMsg()
if received:
break
QtCore.QThread.msleep(5)
self.assertTrue(received)
self.assertLessEqual(
len(received[0]),
CoreProcessWorkerModule.MsgQueue.MAXIMUM_MESSAGE_CHARACTERS,
)
finally:
messageQueue.dispose()
def testCoreLogQueueBacksOffWhileIdleAndRecoversOnActivity(self):
"""Poll rapidly during output and progressively less often while idle."""
received = []
messageQueue = CoreProcessWorkerModule.MsgQueue(msgCallback=received.append)
try:
messageQueue.getNoWait = mock.Mock(return_value='')
expectedTimeouts = (32, 64, 128, 256, 256)
for expectedTimeout in expectedTimeouts:
messageQueue.processMsg()
self.assertEqual(messageQueue.getTimeout(), expectedTimeout)
self.assertEqual(messageQueue.timer.interval(), expectedTimeout)
messageQueue.getNoWait = mock.Mock(side_effect=('message', ''))
messageQueue.processMsg()
self.assertEqual(received, ['message'])
self.assertEqual(
messageQueue.getTimeout(),
CoreProcessWorkerModule.MsgQueue.ACTIVE_DRAIN_INTERVAL,
)
self.assertEqual(
messageQueue.timer.interval(),
CoreProcessWorkerModule.MsgQueue.ACTIVE_DRAIN_INTERVAL,
)
finally:
messageQueue.dispose()
def testFailuresAtMeaningfulStagesRollBackOnlyEarlierStages(self):
expected = {
'storage': ['cleanup storage', 'cleanup plugins'],