Refactor: Implemented the fuller Core refactor around a unified launch spec

Signed-off-by: LorenEteval <loren.eteval@proton.me>
This commit is contained in:
LorenEteval
2026-06-21 15:56:17 +08:00
committed by Loren Eteval
parent 6d28cceacf
commit bcd23cae1e
6 changed files with 306 additions and 119 deletions
+30 -21
View File
@@ -318,6 +318,14 @@ class CoreManager(Mixins.CleanupOnExit):
else:
configcopy = config
def abortStart(message: str = ''):
if message:
logger.error(message)
self.stopAll()
return False
process, success = self._startCore(
configcopy,
routing,
@@ -335,7 +343,9 @@ class CoreManager(Mixins.CleanupOnExit):
if isinstance(process, CoreProcessWorker):
logger.error(f'core {process.name()} start failed')
return success
self.stopAll()
return False
# TUN Mode handling
if not proxyModeOnly and SystemRuntime.isTUNMode():
@@ -376,9 +386,7 @@ class CoreManager(Mixins.CleanupOnExit):
)
if len(defaultGateway) != 1:
logger.error(f'bad default gateway: {defaultGateway}')
return False
return abortStart(f'bad default gateway: {defaultGateway}')
if PLATFORM == 'Windows' or PLATFORM == 'Linux':
# On Linux the 'interface' is a name
@@ -386,9 +394,7 @@ class CoreManager(Mixins.CleanupOnExit):
elif PLATFORM == 'Darwin':
gateway, interface = defaultGateway[0], None
else:
logger.error(f'unrecognized platform: {PLATFORM}')
return False
return abortStart(f'unrecognized platform: {PLATFORM}')
tun = Tun2socks(exitCallback=exitCallback, msgCallback=msgCallbackTUN_)
self.processesPool.append(tun)
@@ -440,7 +446,7 @@ class CoreManager(Mixins.CleanupOnExit):
if PLATFORM != 'Linux':
# Windows & macOS: bring up TUN first
if not startTUN():
return False
return abortStart(f'core {Tun2socks.name()} start failed')
# Handle user defined settings
bypassTUN = userBypassTUNAdapterInterfaceIP()
@@ -457,7 +463,7 @@ class CoreManager(Mixins.CleanupOnExit):
SystemRoutingTable.Relations.clear()
return False
return abortStart()
else:
for bypass in bypassSplit:
if isValidIPAddress(bypass):
@@ -472,7 +478,7 @@ class CoreManager(Mixins.CleanupOnExit):
SystemRoutingTable.Relations.clear()
return False
return abortStart()
else:
logger.info(
f'automatically fetching TUN settings: '
@@ -487,11 +493,9 @@ class CoreManager(Mixins.CleanupOnExit):
error, resolved = DNSResolver.resolve(address)
if error:
logger.error(f'DNS resolution failed: {address}')
SystemRoutingTable.Relations.clear()
return False
return abortStart(f'DNS resolution failed: {address}')
else:
for address in resolved:
SystemRoutingTable.Relations.append([address, gateway])
@@ -504,7 +508,7 @@ class CoreManager(Mixins.CleanupOnExit):
SystemRoutingTable.WIN32IpconfigFindContent,
APPLICATION_TUN_DEVICE_NAME,
):
return False
return abortStart()
# Handle user defined settings
userInterfaceName = userPrimaryAdapterInterfaceName()
@@ -675,17 +679,17 @@ class CoreManager(Mixins.CleanupOnExit):
if not SystemRoutingTable.LinuxExecutePrivilegedScript(
file.name, shell='bash'
):
return False
return abortStart()
if not self.waitForTUNDeviceBroughtUp(
SystemRoutingTable.LinuxFindTUNDevice,
APPLICATION_TUN_DEVICE_NAME,
):
return False
return abortStart()
# Now bring up TUN
if not startTUN():
return False
return abortStart(f'core {Tun2socks.name()} start failed')
return True
@@ -696,11 +700,16 @@ class CoreManager(Mixins.CleanupOnExit):
return any(process.isAlive() for process in self.processesPool)
def stopAll(self):
if self.processesPool:
for process in self.processesPool:
if isinstance(process, CoreProcessFactory):
process.stop()
try:
for process in list(self.processesPool):
if not isinstance(process, CoreProcessFactory):
continue
try:
process.stop()
except Exception as ex:
logger.error(f'error stopping core process: {ex}')
finally:
self.processesPool.clear()
def cleanup(self):
+170 -27
View File
@@ -23,7 +23,9 @@ from Furious.Interface import *
from PySide6 import QtCore
from abc import ABC
from typing import Callable, Union
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Dict, Tuple, Union
import os
import sys
@@ -34,11 +36,73 @@ import threading
import multiprocessing
import multiprocessing.queues
__all__ = ['MsgQueue', 'CoreProcessWorker', 'ProcessOutputRedirector']
__all__ = [
'MsgQueue',
'CoreLaunchSpec',
'CoreProcessState',
'CoreProcessMonitor',
'CoreProcessWorker',
'ProcessOutputRedirector',
]
logger = logging.getLogger(__name__)
class CoreProcessState(Enum):
Idle = 'idle'
Starting = 'starting'
Running = 'running'
Stopping = 'stopping'
Exited = 'exited'
Failed = 'failed'
@dataclass
class CoreLaunchSpec:
target: Callable
args: Tuple[Any, ...] = field(default_factory=tuple)
processKwargs: Dict[str, Any] = field(default_factory=dict)
daemon: bool = True
waitCore: bool = True
waitTime: int = 2500
def __post_init__(self):
self.args = tuple(self.args)
self.processKwargs = dict(self.processKwargs)
self.daemon = self.processKwargs.pop('daemon', self.daemon)
self.waitCore = self.processKwargs.pop('waitCore', self.waitCore)
self.waitTime = self.processKwargs.pop('waitTime', self.waitTime)
@classmethod
def fromProcessKwargs(cls, **kwargs):
daemon = kwargs.pop('daemon', True)
waitCore = kwargs.pop('waitCore', True)
waitTime = kwargs.pop('waitTime', 2500)
target = kwargs.pop('target', None)
args = kwargs.pop('args', tuple())
return cls(
target=target,
args=args,
processKwargs=kwargs,
daemon=daemon,
waitCore=waitCore,
waitTime=waitTime,
)
def toProcessKwargs(self):
kwargs = dict(self.processKwargs)
kwargs.update(
{
'target': self.target,
'args': tuple(self.args),
}
)
return kwargs
class MsgQueue(multiprocessing.queues.Queue):
MSG_PRODUCE_THRESHOLD = 1024
OPTIMIZER_MIN_FREQ = 2
@@ -122,12 +186,16 @@ class MsgQueue(multiprocessing.queues.Queue):
class CoreProcessMonitor(CoreProcessFactory, ABC):
StopJoinTimeout = 3
def __init__(self, **kwargs):
exitCallback = kwargs.pop('exitCallback', None)
super().__init__(exitCallback)
self._process = None
self._state = CoreProcessState.Idle
self._lastExitCode = None
self._daemon = QtCore.QTimer()
self._daemon.timeout.connect(self.queryIsAlive)
@@ -144,6 +212,20 @@ class CoreProcessMonitor(CoreProcessFactory, ABC):
def daemon(self) -> QtCore.QTimer:
return self._daemon
@property
def state(self) -> CoreProcessState:
return self._state
def setState(self, state: CoreProcessState):
self._state = state
@property
def lastExitCode(self):
return self._lastExitCode
def setLastExitCode(self, exitCode):
self._lastExitCode = exitCode
def isAlive(self) -> bool:
if isinstance(self.process, multiprocessing.Process):
return self.process.is_alive()
@@ -155,15 +237,25 @@ class CoreProcessMonitor(CoreProcessFactory, ABC):
if self.process.is_alive():
return True
else:
self.handleInteralProcessStopped()
self.handleInternalProcessStopped()
return False
else:
return False
def handleInteralProcessStopped(self, *args, **kwargs):
def handleInternalProcessStopped(self, *args, **kwargs):
raise NotImplementedError
def closeProcess(self):
if isinstance(self.process, multiprocessing.Process):
try:
self.process.close()
except Exception:
# close() can fail if the process handle is still considered active.
pass
self.process = None
class CoreProcessWorker(CoreProcessMonitor, ABC):
def __init__(self, **kwargs):
@@ -177,37 +269,73 @@ class CoreProcessWorker(CoreProcessMonitor, ABC):
msgCallback=msgCallback, backgroundOptimizer=backgroundOptimizer
)
def handleInteralProcessStopped(self):
logger.error(
f'{self.name()} stopped unexpectedly with exitcode {self.process.exitcode}'
)
def handleInternalProcessStopped(self):
exitcode = self.process.exitcode
logger.error(f'{self.name()} stopped unexpectedly with exitcode {exitcode}')
self.msgQueue.stopTimer()
self.daemon.stop()
self.callExitCallback(self.process.exitcode)
self.setLastExitCode(exitcode)
self.setState(CoreProcessState.Failed)
self.callExitCallback(exitcode)
# Reset internal process
self.process = None
self.closeProcess()
def start(self, **kwargs) -> bool:
daemon = kwargs.pop('daemon', True)
waitCore = kwargs.pop('waitCore', True)
waitTime = kwargs.pop('waitTime', 2500)
return self.startWithSpec(CoreLaunchSpec.fromProcessKwargs(**kwargs))
self.process = multiprocessing.Process(**kwargs, daemon=daemon)
self.process.start()
def startWithSpec(self, launchSpec: CoreLaunchSpec) -> bool:
if not isinstance(launchSpec, CoreLaunchSpec):
logger.error(f'invalid launch spec for {self.name()}: {launchSpec}')
self.setState(CoreProcessState.Failed)
return False
if not callable(launchSpec.target):
logger.error(
f'invalid launch target for {self.name()}: {launchSpec.target}'
)
self.setState(CoreProcessState.Failed)
return False
if self.isAlive():
logger.warning(f'{self.name()} is already running. Stop it before restart')
self.stop()
self.setState(CoreProcessState.Starting)
try:
self.process = multiprocessing.Process(
**launchSpec.toProcessKwargs(), daemon=launchSpec.daemon
)
self.process.start()
except Exception as ex:
logger.error(f'{self.name()} start failed: {ex}')
self.setState(CoreProcessState.Failed)
self.closeProcess()
return False
logger.info(f'{self.name()} {self.version()} started')
self.msgQueue.setTimeout(MsgQueue.MSG_PRODUCE_THRESHOLD)
self.msgQueue.startTimer()
if waitCore:
if launchSpec.waitCore:
# Wait for the core to start up completely
PySide6Legacy.eventLoopWait(waitTime)
PySide6Legacy.eventLoopWait(launchSpec.waitTime)
if self.queryIsAlive():
self.setState(CoreProcessState.Running)
# Start core daemon
self.daemon.start(CORE_CHECK_ALIVE_INTERVAL)
@@ -216,19 +344,34 @@ class CoreProcessWorker(CoreProcessMonitor, ABC):
return False
def stop(self):
if self.isAlive():
self.msgQueue.stopTimer()
self.daemon.stop()
self.msgQueue.stopTimer()
self.daemon.stop()
if not isinstance(self.process, multiprocessing.Process):
return
self.setState(CoreProcessState.Stopping)
if self.process.is_alive():
self.process.terminate()
self.process.join()
self.process.join(CoreProcessMonitor.StopJoinTimeout)
logger.info(
f'{self.name()} terminated with exitcode {self.process.exitcode}'
)
else:
# Do nothing
pass
if self.process.is_alive():
logger.warning(
f'{self.name()} did not terminate in '
f'{CoreProcessMonitor.StopJoinTimeout}s. Kill it'
)
self.process.kill()
self.process.join(CoreProcessMonitor.StopJoinTimeout)
exitcode = self.process.exitcode
logger.info(f'{self.name()} terminated with exitcode {exitcode}')
self.setLastExitCode(exitcode)
self.setState(CoreProcessState.Exited)
self.closeProcess()
class ProcessOutputRedirector:
+35 -46
View File
@@ -64,16 +64,16 @@ class Hysteria1(CoreProcessWorker):
super().__init__(**kwargs)
@staticmethod
def rule(rulePath):
if isinstance(rulePath, str) and rulePath == '':
def loadOptionalFile(pathLike, fileType: str):
if isinstance(pathLike, str) and pathLike == '':
return ''
try:
path = absolutePath(str(rulePath))
path = absolutePath(str(pathLike))
except Exception:
# Any non-exit exceptions
logger.error('invalid hysteria1 rule path. Fall back to empty')
logger.error(f'invalid hysteria1 {fileType} path. Fall back to empty')
return ''
@@ -81,47 +81,26 @@ class Hysteria1(CoreProcessWorker):
with open(path, 'rb') as file:
data = file.read()
logger.info(f'hysteria1 rule \'{path}\' load success')
logger.info(f'hysteria1 {fileType} \'{path}\' load success')
return data
except Exception as ex:
# Any non-exit exceptions
logger.error(
f'hysteria1 rule \'{path}\' load failed. {ex}. Fall back to empty'
f'hysteria1 {fileType} \'{path}\' load failed. {ex}. '
f'Fall back to empty'
)
return ''
@staticmethod
def rule(rulePath):
return Hysteria1.loadOptionalFile(rulePath, 'rule')
@staticmethod
def mmdb(mmdbPath):
if isinstance(mmdbPath, str) and mmdbPath == '':
return ''
try:
path = absolutePath(str(mmdbPath))
except Exception:
# Any non-exit exceptions
logger.error('invalid hysteria1 mmdb path. Fall back to empty')
return ''
try:
with open(path, 'rb') as file:
data = file.read()
logger.info(f'hysteria1 mmdb \'{path}\' load success')
return data
except Exception as ex:
# Any non-exit exceptions
logger.error(
f'hysteria1 mmdb \'{path}\' load failed. {ex}. Fall back to empty'
)
return ''
return Hysteria1.loadOptionalFile(mmdbPath, 'mmdb')
@staticmethod
def name() -> str:
@@ -138,19 +117,29 @@ class Hysteria1(CoreProcessWorker):
return '0.0.0'
def start(self, config: Union[str, dict], rule, mmdb, **kwargs) -> bool:
def launchSpec(
self, config: Union[str, dict], rule, mmdb, **kwargs
) -> Union[CoreLaunchSpec, None]:
param = self.toJSONString(config)
if param:
return super().start(
target=startHysteria1,
args=(
param,
rule,
mmdb,
self.msgQueue,
),
**kwargs,
)
else:
if not param:
return None
return CoreLaunchSpec(
target=startHysteria1,
args=(
param,
rule,
mmdb,
self.msgQueue,
),
processKwargs=kwargs,
)
def start(self, config: Union[str, dict], rule, mmdb, **kwargs) -> bool:
launchSpec = self.launchSpec(config, rule, mmdb, **kwargs)
if launchSpec is None:
return False
return self.startWithSpec(launchSpec)
+21 -11
View File
@@ -79,17 +79,27 @@ class Hysteria2(CoreProcessWorker):
return '0.0.0'
def start(self, config: Union[str, dict], **kwargs) -> bool:
def launchSpec(
self, config: Union[str, dict], **kwargs
) -> Union[CoreLaunchSpec, None]:
param = self.toJSONString(config)
if param:
return super().start(
target=startHysteria2,
args=(
param,
self.msgQueue,
),
**kwargs,
)
else:
if not param:
return None
return CoreLaunchSpec(
target=startHysteria2,
args=(
param,
self.msgQueue,
),
processKwargs=kwargs,
)
def start(self, config: Union[str, dict], **kwargs) -> bool:
launchSpec = self.launchSpec(config, **kwargs)
if launchSpec is None:
return False
return self.startWithSpec(launchSpec)
+29 -3
View File
@@ -75,7 +75,7 @@ class Tun2socks(CoreProcessWorker):
return '0.0.0'
def start(
def launchSpec(
self,
device: str,
networkInterface: str,
@@ -86,8 +86,8 @@ class Tun2socks(CoreProcessWorker):
tcpReceiveBufferSize: str = '',
tcpAutoTuning: bool = False,
**kwargs,
) -> bool:
return super().start(
) -> CoreLaunchSpec:
return CoreLaunchSpec(
target=startTun2socks,
args=(
self.msgQueue,
@@ -100,9 +100,35 @@ class Tun2socks(CoreProcessWorker):
tcpReceiveBufferSize,
tcpAutoTuning,
),
processKwargs=kwargs,
)
def start(
self,
device: str,
networkInterface: str,
logLevel: str,
proxy: str,
restAPI: str,
tcpSendBufferSize: str = '',
tcpReceiveBufferSize: str = '',
tcpAutoTuning: bool = False,
**kwargs,
) -> bool:
launchSpec = self.launchSpec(
device,
networkInterface,
logLevel,
proxy,
restAPI,
tcpSendBufferSize,
tcpReceiveBufferSize,
tcpAutoTuning,
**kwargs,
)
return self.startWithSpec(launchSpec)
def stop(self):
if PLATFORM == 'Linux':
# Stop tunnel first
+21 -11
View File
@@ -141,17 +141,27 @@ class XrayCore(CoreProcessWorker):
return '0.0.0'
def start(self, config: Union[str, ConfigFactory, dict], **kwargs) -> bool:
def launchSpec(
self, config: Union[str, ConfigFactory, dict], **kwargs
) -> Union[CoreLaunchSpec, None]:
param = self.toJSONString(config)
if param:
return super().start(
target=startXrayCore,
args=(
param,
self.msgQueue,
),
**kwargs,
)
else:
if not param:
return None
return CoreLaunchSpec(
target=startXrayCore,
args=(
param,
self.msgQueue,
),
processKwargs=kwargs,
)
def start(self, config: Union[str, ConfigFactory, dict], **kwargs) -> bool:
launchSpec = self.launchSpec(config, **kwargs)
if launchSpec is None:
return False
return self.startWithSpec(launchSpec)