mirror of
https://github.com/LorenEteval/Furious.git
synced 2026-10-08 22:59:48 +03:00
Rename managed core lifecycle to CoreRuntime
Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
@@ -15,15 +15,15 @@
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
"""Integrate user-managed local executables with the Furious kernel API."""
|
||||
"""Integrate user-managed local executables with the core-runtime API."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from Furious.Plugins.API import (
|
||||
FuriousPlugin,
|
||||
KernelFactory,
|
||||
KernelLaunch,
|
||||
KernelRequest,
|
||||
CoreRuntimeFactory,
|
||||
CoreRuntimeLaunch,
|
||||
CoreRuntimeRequest,
|
||||
PluginMetadata,
|
||||
)
|
||||
|
||||
@@ -39,20 +39,20 @@ __all__ = ['ExternalCorePlugin']
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ExternalCoreKernelFactory(KernelFactory):
|
||||
class ExternalCoreRuntimeFactory(CoreRuntimeFactory):
|
||||
"""Construct a managed direct-process runtime from a structured profile."""
|
||||
|
||||
factoryId = 'official.external-core'
|
||||
configurationTypes = (ConfigExternalCore,)
|
||||
kernelTypes = (ExternalCoreProcess,)
|
||||
runtimeTypes = (ExternalCoreProcess,)
|
||||
|
||||
def usesApplicationTun2socks(self, config) -> bool:
|
||||
"""Honor this profile's explicit host-managed tun2socks preference."""
|
||||
return config.usesApplicationTun2socks()
|
||||
|
||||
def create(self, request: KernelRequest):
|
||||
def create(self, request: CoreRuntimeRequest):
|
||||
"""Create an External Core process launch for the connection manager."""
|
||||
process = ExternalCoreProcess(
|
||||
runtime = ExternalCoreProcess(
|
||||
exitCallback=request.exitCallback,
|
||||
msgCallback=request.messageCallback,
|
||||
)
|
||||
@@ -72,8 +72,8 @@ class ExternalCoreKernelFactory(KernelFactory):
|
||||
'application-managed TUN will be skipped'
|
||||
)
|
||||
|
||||
return KernelLaunch(
|
||||
process,
|
||||
return CoreRuntimeLaunch(
|
||||
runtime,
|
||||
request.configuration,
|
||||
options=request.options,
|
||||
)
|
||||
@@ -98,5 +98,5 @@ class ExternalCorePlugin(FuriousPlugin):
|
||||
self.capabilities = (
|
||||
*EXTERNAL_CORE_PROTOCOL_HANDLERS,
|
||||
*EXTERNAL_CORE_PROTOCOL_EDITORS,
|
||||
ExternalCoreKernelFactory(),
|
||||
ExternalCoreRuntimeFactory(),
|
||||
)
|
||||
|
||||
@@ -20,7 +20,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from Furious.Core import CoreProcessState
|
||||
from Furious.Interface import RuntimeKernel
|
||||
from Furious.Interface import CoreRuntime
|
||||
|
||||
from .Configuration import ConfigExternalCore
|
||||
|
||||
@@ -38,7 +38,7 @@ __all__ = ['ExternalCoreProcess']
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ExternalCoreProcess(RuntimeKernel):
|
||||
class ExternalCoreProcess(CoreRuntime):
|
||||
"""Manage one direct external process and its output-reader threads."""
|
||||
|
||||
StartupObservationTimeout = 0.25
|
||||
|
||||
@@ -34,12 +34,12 @@ __all__ = ['Hysteria1Plugin']
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Hysteria1KernelFactory(KernelFactory):
|
||||
"""Construct Hysteria 1 kernels independently of protocol handling."""
|
||||
class Hysteria1CoreRuntimeFactory(CoreRuntimeFactory):
|
||||
"""Construct Hysteria 1 core runtimes independently of protocol handling."""
|
||||
|
||||
factoryId = 'official.hysteria1'
|
||||
configurationTypes = (ConfigHysteria1,)
|
||||
kernelTypes = (Hysteria1,)
|
||||
runtimeTypes = (Hysteria1,)
|
||||
|
||||
def routingOptions(self, config=None):
|
||||
"""Return the routing modes supported by Hysteria 1."""
|
||||
@@ -48,8 +48,8 @@ class Hysteria1KernelFactory(KernelFactory):
|
||||
for routing in AppBuiltinRouting
|
||||
)
|
||||
|
||||
def create(self, request: KernelRequest):
|
||||
"""Configure routing and create a Hysteria 1 kernel launch."""
|
||||
def create(self, request: CoreRuntimeRequest):
|
||||
"""Configure routing and create a Hysteria 1 core-runtime launch."""
|
||||
config, routing = (
|
||||
request.configuration,
|
||||
request.routing,
|
||||
@@ -81,13 +81,13 @@ class Hysteria1KernelFactory(KernelFactory):
|
||||
logger.info(f'routing is {routing}')
|
||||
logger.info(f'RoutingObject: {routingObject}')
|
||||
|
||||
process = Hysteria1(
|
||||
runtime = Hysteria1(
|
||||
exitCallback=request.exitCallback,
|
||||
msgCallback=request.messageCallback,
|
||||
)
|
||||
|
||||
return KernelLaunch(
|
||||
process,
|
||||
return CoreRuntimeLaunch(
|
||||
runtime,
|
||||
config,
|
||||
arguments=(
|
||||
Hysteria1.rule(routingObject.get('rule', '')),
|
||||
@@ -135,5 +135,5 @@ class Hysteria1Plugin(FuriousPlugin):
|
||||
self.capabilities = (
|
||||
*HYSTERIA1_PROTOCOL_HANDLERS,
|
||||
*HYSTERIA1_PROTOCOL_EDITORS,
|
||||
Hysteria1KernelFactory(),
|
||||
Hysteria1CoreRuntimeFactory(),
|
||||
)
|
||||
|
||||
@@ -153,12 +153,12 @@ class Hysteria2ActionProvider(ActionProvider):
|
||||
)
|
||||
|
||||
|
||||
class Hysteria2KernelFactory(KernelFactory):
|
||||
"""Construct Hysteria 2 kernels independently of protocol handling."""
|
||||
class Hysteria2CoreRuntimeFactory(CoreRuntimeFactory):
|
||||
"""Construct Hysteria 2 core runtimes independently of protocol handling."""
|
||||
|
||||
factoryId = 'official.hysteria2'
|
||||
configurationTypes = (ConfigHysteria2,)
|
||||
kernelTypes = (Hysteria2,)
|
||||
runtimeTypes = (Hysteria2,)
|
||||
|
||||
def prepareTUN(self, config) -> bool:
|
||||
"""Preserve user TUN or replace it with Furious-managed native TUN."""
|
||||
@@ -195,20 +195,20 @@ class Hysteria2KernelFactory(KernelFactory):
|
||||
"""Use host tun2socks only when the runtime has no native TUN block."""
|
||||
return not hasHysteria2TUNConfig(config)
|
||||
|
||||
def create(self, request: KernelRequest):
|
||||
"""Create a prepared Hysteria 2 kernel launch."""
|
||||
def create(self, request: CoreRuntimeRequest):
|
||||
"""Create a prepared Hysteria 2 core-runtime launch."""
|
||||
if request.log:
|
||||
logger.info(f'core {Hysteria2.name()} configured')
|
||||
|
||||
process = Hysteria2(
|
||||
runtime = Hysteria2(
|
||||
exitCallback=request.exitCallback,
|
||||
msgCallback=request.messageCallback,
|
||||
)
|
||||
|
||||
setattr(process, 'hysteria2StatsTarget', configuredHysteria2StatsTarget())
|
||||
setattr(runtime, 'hysteria2StatsTarget', configuredHysteria2StatsTarget())
|
||||
|
||||
return KernelLaunch(
|
||||
process,
|
||||
return CoreRuntimeLaunch(
|
||||
runtime,
|
||||
request.configuration,
|
||||
options=request.options,
|
||||
)
|
||||
@@ -261,7 +261,7 @@ class Hysteria2Plugin(FuriousPlugin):
|
||||
self.capabilities = (
|
||||
*HYSTERIA2_PROTOCOL_HANDLERS,
|
||||
*HYSTERIA2_PROTOCOL_EDITORS,
|
||||
Hysteria2KernelFactory(),
|
||||
Hysteria2CoreRuntimeFactory(),
|
||||
Hysteria2StatsProvider(),
|
||||
Hysteria2SettingsProvider(),
|
||||
Hysteria2ActionProvider(),
|
||||
|
||||
@@ -196,11 +196,11 @@ class Hysteria2StatsProvider(TrafficStatsProvider):
|
||||
"""Expose Hysteria server counters through the shared metrics capability."""
|
||||
|
||||
providerId = 'official.hysteria2.stats'
|
||||
kernelTypes = (Hysteria2,)
|
||||
runtimeTypes = (Hysteria2,)
|
||||
|
||||
def monitorForKernel(self, kernel) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor when the active kernel has a configured API target."""
|
||||
target = getattr(kernel, 'hysteria2StatsTarget', None)
|
||||
def monitorForRuntime(self, runtime) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor when the active runtime has a configured API target."""
|
||||
target = getattr(runtime, 'hysteria2StatsTarget', None)
|
||||
|
||||
if not isinstance(target, Hysteria2StatsTarget):
|
||||
return None
|
||||
|
||||
@@ -152,12 +152,12 @@ class XrayActionProvider(ActionProvider):
|
||||
)
|
||||
|
||||
|
||||
class XrayKernelFactory(KernelFactory):
|
||||
"""Construct Xray kernels independently of protocols and editors."""
|
||||
class XrayCoreRuntimeFactory(CoreRuntimeFactory):
|
||||
"""Construct Xray core runtimes independently of protocols and editors."""
|
||||
|
||||
factoryId = 'official.xray'
|
||||
configurationTypes = (ConfigXray,)
|
||||
kernelTypes = (XrayCore,)
|
||||
runtimeTypes = (XrayCore,)
|
||||
|
||||
def fromMapping(self, configuration, **kwargs):
|
||||
"""Recognize a complete Xray configuration without a proxy profile."""
|
||||
@@ -237,8 +237,8 @@ class XrayKernelFactory(KernelFactory):
|
||||
"""Point Xray-core at Furious's bundled geo-asset directory."""
|
||||
os.environ['XRAY_LOCATION_ASSET'] = str(XRAY_ASSET_DIR)
|
||||
|
||||
def create(self, request: KernelRequest):
|
||||
"""Configure routing and create an Xray-core launch."""
|
||||
def create(self, request: CoreRuntimeRequest):
|
||||
"""Configure routing and create an Xray core-runtime launch."""
|
||||
config, routing, proxyModeOnly, log = (
|
||||
request.configuration,
|
||||
request.routing,
|
||||
@@ -321,13 +321,13 @@ class XrayKernelFactory(KernelFactory):
|
||||
|
||||
statsTarget = None
|
||||
|
||||
process = XrayCore(
|
||||
runtime = XrayCore(
|
||||
exitCallback=request.exitCallback,
|
||||
msgCallback=request.messageCallback,
|
||||
)
|
||||
process.xrayStatsTarget = statsTarget
|
||||
runtime.xrayStatsTarget = statsTarget
|
||||
|
||||
return KernelLaunch(process, config, options=request.options)
|
||||
return CoreRuntimeLaunch(runtime, config, options=request.options)
|
||||
|
||||
def prepareDownloadTest(self, config, port: int):
|
||||
"""Create an Xray configuration with one local HTTP test inbound."""
|
||||
@@ -426,7 +426,7 @@ class XrayPlugin(FuriousPlugin):
|
||||
self.capabilities = (
|
||||
*XRAY_PROTOCOL_HANDLERS,
|
||||
*XRAY_PROTOCOL_EDITORS,
|
||||
XrayKernelFactory(),
|
||||
XrayCoreRuntimeFactory(),
|
||||
XrayStatsProvider(),
|
||||
XrayActionProvider(),
|
||||
)
|
||||
|
||||
@@ -271,11 +271,11 @@ class XrayStatsProvider(TrafficStatsProvider):
|
||||
"""Expose active Xray outbound traffic through the plugin capability API."""
|
||||
|
||||
providerId = 'official.xray.stats'
|
||||
kernelTypes = (XrayCore,)
|
||||
runtimeTypes = (XrayCore,)
|
||||
|
||||
def monitorForKernel(self, kernel) -> Optional[TrafficStatsMonitor]:
|
||||
def monitorForRuntime(self, runtime) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor when the active core has a valid outbound target."""
|
||||
target = getattr(kernel, 'xrayStatsTarget', None)
|
||||
target = getattr(runtime, 'xrayStatsTarget', None)
|
||||
|
||||
if not isinstance(target, XrayStatsTarget):
|
||||
return None
|
||||
|
||||
@@ -86,7 +86,7 @@ class ConnectionController(QtCore.QObject):
|
||||
stateChanged = QtCore.Signal(object)
|
||||
activeConfigurationChanged = QtCore.Signal(object)
|
||||
interactionEnabledChanged = QtCore.Signal(bool)
|
||||
processesChanged = QtCore.Signal(object)
|
||||
runtimesChanged = QtCore.Signal(object)
|
||||
progressStarted = QtCore.Signal()
|
||||
progressFinished = QtCore.Signal(bool)
|
||||
notificationRequested = QtCore.Signal(str)
|
||||
@@ -117,9 +117,9 @@ class ConnectionController(QtCore.QObject):
|
||||
return self._activeConfiguration
|
||||
|
||||
@property
|
||||
def processes(self):
|
||||
"""Return an immutable snapshot of managed core processes."""
|
||||
return tuple(self._coreManager.processesPool)
|
||||
def runtimes(self):
|
||||
"""Return an immutable snapshot of managed core runtimes."""
|
||||
return tuple(self._coreManager.runtimes)
|
||||
|
||||
@property
|
||||
def lastError(self):
|
||||
@@ -167,9 +167,9 @@ class ConnectionController(QtCore.QObject):
|
||||
self._activeConfiguration = configuration
|
||||
self.activeConfigurationChanged.emit(configuration)
|
||||
|
||||
def _setProcessesChanged(self):
|
||||
"""Publish a stable process snapshot after runtime changes."""
|
||||
self.processesChanged.emit(self.processes)
|
||||
def _emitRuntimesChanged(self):
|
||||
"""Publish a stable snapshot after managed runtimes change."""
|
||||
self.runtimesChanged.emit(self.runtimes)
|
||||
|
||||
def _reportError(self, message: str, details: str = ''):
|
||||
"""Publish a user-facing error without choosing its presentation."""
|
||||
@@ -296,7 +296,7 @@ class ConnectionController(QtCore.QObject):
|
||||
success = False
|
||||
startExceptionDetails = str(ex)
|
||||
|
||||
self._setProcessesChanged()
|
||||
self._emitRuntimesChanged()
|
||||
|
||||
if not self._actionQueue.empty():
|
||||
while not self._actionQueue.empty():
|
||||
@@ -372,7 +372,7 @@ class ConnectionController(QtCore.QObject):
|
||||
# strand every connection UI in the disabled Disconnecting state.
|
||||
logger.error(f'failed to stop connection runtime: {ex}')
|
||||
|
||||
self._setProcessesChanged()
|
||||
self._emitRuntimesChanged()
|
||||
self._reset()
|
||||
|
||||
while not self._actionQueue.empty():
|
||||
@@ -463,7 +463,7 @@ class ConnectionController(QtCore.QObject):
|
||||
if callable(action):
|
||||
action()
|
||||
|
||||
def coreExitCallback(self, core: CoreProcess, exitcode: int):
|
||||
def coreExitCallback(self, core: CoreRuntime, exitcode: int):
|
||||
"""Translate a core exit into a queued lifecycle operation."""
|
||||
|
||||
def putItem(item):
|
||||
@@ -475,12 +475,12 @@ class ConnectionController(QtCore.QObject):
|
||||
|
||||
pass
|
||||
|
||||
if exitcode == CoreProcess.ExitCode.SystemShuttingDown.value:
|
||||
if exitcode == CoreRuntime.ExitCode.SystemShuttingDown.value:
|
||||
return None
|
||||
|
||||
if exitcode == CoreProcess.ExitCode.ConfigurationError.value:
|
||||
if exitcode == CoreRuntime.ExitCode.ConfigurationError.value:
|
||||
message = f'{core.name()}: ' + _('Invalid server configuration')
|
||||
elif exitcode == CoreProcess.ExitCode.ServerStartFailure.value:
|
||||
elif exitcode == CoreRuntime.ExitCode.ServerStartFailure.value:
|
||||
message = f'{core.name()}: ' + _('Failed to start core')
|
||||
else:
|
||||
pluginMessage = getPluginRegistry().coreExitMessage(core, exitcode)
|
||||
|
||||
@@ -226,7 +226,7 @@ class MsgQueue(multiprocessing.queues.Queue):
|
||||
self.startTimer()
|
||||
|
||||
|
||||
class CoreProcessMonitor(CoreProcess, ABC):
|
||||
class CoreProcessMonitor(CoreRuntime, ABC):
|
||||
"""Track the state and lifetime of a proxy-core child process."""
|
||||
|
||||
StopJoinTimeout = 3
|
||||
@@ -312,7 +312,7 @@ class CoreProcessMonitor(CoreProcess, ABC):
|
||||
self.process = None
|
||||
|
||||
def dispose(self):
|
||||
"""Release the monitor timer after this kernel leaves its owner pool."""
|
||||
"""Release the monitor timer after this runtime leaves its owner pool."""
|
||||
self.daemon.stop()
|
||||
|
||||
try:
|
||||
|
||||
@@ -15,7 +15,7 @@
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
"""Define the common proxy-core process interface."""
|
||||
"""Define the lifecycle contract for one managed proxy-core runtime."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -26,21 +26,27 @@ from abc import ABC, abstractmethod
|
||||
from enum import Enum
|
||||
from typing import Callable, Union
|
||||
|
||||
import logging
|
||||
import functools
|
||||
import logging
|
||||
|
||||
__all__ = ['CoreProcess', 'RuntimeKernel']
|
||||
__all__ = ['CoreRuntime']
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_INVALID_CONFIGURATION_ERROR = 'Invalid server configuration'
|
||||
|
||||
|
||||
class RuntimeKernel(ABC):
|
||||
"""Define the lifecycle shared by plugin-created runtime kernels."""
|
||||
class CoreRuntime(ABC):
|
||||
"""Manage one startable and stoppable proxy-core runtime instance.
|
||||
|
||||
A runtime represents the semantic lifecycle of one proxy core. Concrete
|
||||
implementations may use a subprocess, multiprocessing, or an in-process
|
||||
binding; those execution mechanisms are intentionally outside this
|
||||
contract.
|
||||
"""
|
||||
|
||||
class ExitCode(Enum):
|
||||
"""Enumerate process exit codes."""
|
||||
"""Enumerate shared semantic core exit codes."""
|
||||
|
||||
ConfigurationError = 23
|
||||
# Windows: 4294967295. Darwin, Linux: 255 (-1)
|
||||
@@ -50,16 +56,16 @@ class RuntimeKernel(ABC):
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
exitCallback: Union[Callable[[RuntimeKernel, int], None], None] = None,
|
||||
exitCallback: Union[Callable[[CoreRuntime, int], None], None] = None,
|
||||
):
|
||||
"""Initialize the core process."""
|
||||
"""Initialize an idle core runtime."""
|
||||
super().__init__()
|
||||
|
||||
self._exitCallback = exitCallback
|
||||
self._startError = ''
|
||||
|
||||
def callExitCallback(self, exitcode: int):
|
||||
"""Call exit callback."""
|
||||
"""Report this runtime's termination to its lifecycle owner."""
|
||||
if callable(self._exitCallback):
|
||||
self._exitCallback(self, exitcode)
|
||||
|
||||
@@ -77,7 +83,7 @@ class RuntimeKernel(ABC):
|
||||
|
||||
@functools.singledispatchmethod
|
||||
def toJSONString(self, config, **kwargs) -> str:
|
||||
"""Serialize the configuration as JSON text."""
|
||||
"""Serialize a prepared runtime configuration as JSON text."""
|
||||
self.setStartError(_INVALID_CONFIGURATION_ERROR)
|
||||
|
||||
logger.error(
|
||||
@@ -89,7 +95,7 @@ class RuntimeKernel(ABC):
|
||||
|
||||
@toJSONString.register(str)
|
||||
def _(self, config, **kwargs) -> str:
|
||||
"""Handle the registered singledispatch variant."""
|
||||
"""Accept an already serialized runtime configuration."""
|
||||
if not config:
|
||||
self.setStartError(_INVALID_CONFIGURATION_ERROR)
|
||||
|
||||
@@ -103,7 +109,7 @@ class RuntimeKernel(ABC):
|
||||
|
||||
@toJSONString.register(dict)
|
||||
def _(self, config, **kwargs) -> str:
|
||||
"""Handle the registered singledispatch variant."""
|
||||
"""Serialize a mapping with its own serializer or the shared encoder."""
|
||||
serializer = getattr(config, 'toJSONString', None)
|
||||
|
||||
try:
|
||||
@@ -142,25 +148,21 @@ class RuntimeKernel(ABC):
|
||||
@staticmethod
|
||||
@abstractmethod
|
||||
def name() -> str:
|
||||
"""Return the process implementation name."""
|
||||
"""Return the proxy-core implementation name."""
|
||||
raise NotImplementedError
|
||||
|
||||
@staticmethod
|
||||
@abstractmethod
|
||||
def version() -> str:
|
||||
"""Return the bundled core version."""
|
||||
"""Return the bundled core version, if one exists."""
|
||||
raise NotImplementedError
|
||||
|
||||
@abstractmethod
|
||||
def start(self, *args, **kwargs) -> bool:
|
||||
"""Start the runtime kernel."""
|
||||
"""Start this core runtime."""
|
||||
raise NotImplementedError
|
||||
|
||||
@abstractmethod
|
||||
def stop(self):
|
||||
"""Stop the runtime kernel."""
|
||||
"""Stop this core runtime."""
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
# Compatibility name retained for existing process implementations.
|
||||
CoreProcess = RuntimeKernel
|
||||
@@ -21,14 +21,13 @@ from __future__ import annotations
|
||||
|
||||
from .Application import ApplicationRunner
|
||||
from .Editor import EditorBinding, EditorWidgetBinding
|
||||
from .Process import CoreProcess, RuntimeKernel
|
||||
from .Runtime import CoreRuntime
|
||||
from .Storage import StorageBackend
|
||||
|
||||
__all__ = [
|
||||
'ApplicationRunner',
|
||||
'CoreProcess',
|
||||
'CoreRuntime',
|
||||
'EditorBinding',
|
||||
'EditorWidgetBinding',
|
||||
'StorageBackend',
|
||||
'RuntimeKernel',
|
||||
]
|
||||
|
||||
+22
-22
@@ -27,10 +27,10 @@ __all__ = [
|
||||
'PLUGIN_API_VERSION',
|
||||
'ActionProvider',
|
||||
'CapabilityKind',
|
||||
'CoreRuntimeFactory',
|
||||
'CoreRuntimeLaunch',
|
||||
'CoreRuntimeRequest',
|
||||
'FuriousPlugin',
|
||||
'KernelFactory',
|
||||
'KernelLaunch',
|
||||
'KernelRequest',
|
||||
'NavigationPageDescriptor',
|
||||
'NavigationPageProvider',
|
||||
'PluginCapability',
|
||||
@@ -68,7 +68,7 @@ class CapabilityKind(str, Enum):
|
||||
Protocol = 'protocol'
|
||||
ProtocolEditor = 'protocol-editor'
|
||||
SubscriptionDecoder = 'subscription-decoder'
|
||||
KernelFactory = 'kernel-factory'
|
||||
CoreRuntimeFactory = 'core-runtime-factory'
|
||||
TrafficStats = 'traffic-stats'
|
||||
PluginSettings = 'plugin-settings'
|
||||
NavigationPage = 'navigation-page'
|
||||
@@ -371,25 +371,25 @@ class TrafficStatsMonitor:
|
||||
|
||||
|
||||
class TrafficStatsProvider(PluginCapability):
|
||||
"""Provide traffic counters for one or more runtime kernel types."""
|
||||
"""Provide traffic counters for one or more core-runtime types."""
|
||||
|
||||
capabilityKind = CapabilityKind.TrafficStats
|
||||
providerId = ''
|
||||
kernelTypes = tuple()
|
||||
runtimeTypes = tuple()
|
||||
|
||||
@property
|
||||
def capabilityId(self) -> str:
|
||||
"""Return the traffic-statistics provider identifier."""
|
||||
return self.providerId
|
||||
|
||||
def monitorForKernel(self, kernel) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor for *kernel* or ``None`` when unavailable."""
|
||||
def monitorForRuntime(self, runtime) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor for *runtime* or ``None`` when unavailable."""
|
||||
return None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KernelRequest:
|
||||
"""Describe one runtime-kernel construction request."""
|
||||
class CoreRuntimeRequest:
|
||||
"""Describe one core-runtime construction request."""
|
||||
|
||||
configuration: Any
|
||||
routing: str
|
||||
@@ -401,18 +401,18 @@ class KernelRequest:
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KernelLaunch:
|
||||
"""Bind a constructed kernel to its prepared start arguments."""
|
||||
class CoreRuntimeLaunch:
|
||||
"""Bind a constructed core runtime to its prepared start arguments."""
|
||||
|
||||
kernel: Any
|
||||
runtime: Any
|
||||
configuration: Any
|
||||
arguments: Tuple[Any, ...] = tuple()
|
||||
options: Mapping[str, Any] = field(default_factory=dict)
|
||||
|
||||
def start(self) -> bool:
|
||||
"""Start the prepared kernel."""
|
||||
"""Start the prepared core runtime."""
|
||||
return bool(
|
||||
self.kernel.start(
|
||||
self.runtime.start(
|
||||
self.configuration,
|
||||
*self.arguments,
|
||||
**dict(self.options),
|
||||
@@ -420,13 +420,13 @@ class KernelLaunch:
|
||||
)
|
||||
|
||||
|
||||
class KernelFactory(PluginCapability):
|
||||
"""Construct runtime kernels independently from protocol handling."""
|
||||
class CoreRuntimeFactory(PluginCapability):
|
||||
"""Construct managed core runtimes independently from protocol handling."""
|
||||
|
||||
capabilityKind = CapabilityKind.KernelFactory
|
||||
capabilityKind = CapabilityKind.CoreRuntimeFactory
|
||||
factoryId = ''
|
||||
configurationTypes = tuple()
|
||||
kernelTypes = tuple()
|
||||
runtimeTypes = tuple()
|
||||
|
||||
@property
|
||||
def capabilityId(self) -> str:
|
||||
@@ -457,10 +457,10 @@ class KernelFactory(PluginCapability):
|
||||
return tuple()
|
||||
|
||||
def configureEnvironment(self):
|
||||
"""Set optional environment required by this backend's process."""
|
||||
"""Set optional environment required by this backend's runtime."""
|
||||
|
||||
def create(self, request: KernelRequest) -> Optional[KernelLaunch]:
|
||||
"""Create a prepared runtime kernel launch."""
|
||||
def create(self, request: CoreRuntimeRequest) -> Optional[CoreRuntimeLaunch]:
|
||||
"""Create a prepared core-runtime launch."""
|
||||
return None
|
||||
|
||||
def prepareDownloadTest(self, config, port: int):
|
||||
|
||||
+96
-80
@@ -53,6 +53,11 @@ def _normalizeScheme(value) -> str:
|
||||
return str(value).strip().rstrip(':').casefold()
|
||||
|
||||
|
||||
def _runtimeTypes(capability) -> tuple:
|
||||
"""Return the capability's declared managed runtime types."""
|
||||
return tuple(getattr(capability, 'runtimeTypes', tuple()) or tuple())
|
||||
|
||||
|
||||
def _schemeFromURI(uri: str) -> str:
|
||||
"""Extract a normalized scheme from *uri*."""
|
||||
try:
|
||||
@@ -84,7 +89,7 @@ class PluginRegistry:
|
||||
self._protocolEditors = {}
|
||||
self._factories = {}
|
||||
self._configurationFactories = {}
|
||||
self._kernelFactories = {}
|
||||
self._runtimeFactories = {}
|
||||
self._trafficStatsProviders = {}
|
||||
self._decoders = {}
|
||||
self._initializedPlugins = []
|
||||
@@ -150,8 +155,8 @@ class PluginRegistry:
|
||||
localSchemes = set()
|
||||
localEditorProtocols = set()
|
||||
localConfigurationTypes = []
|
||||
localKernelTypes = []
|
||||
localTrafficStatsKernelTypes = []
|
||||
localRuntimeTypes = []
|
||||
localTrafficStatsRuntimeTypes = []
|
||||
entries = []
|
||||
|
||||
for capability in capabilities:
|
||||
@@ -244,13 +249,13 @@ class PluginRegistry:
|
||||
|
||||
localEditorProtocols.update(protocolIds)
|
||||
detail = protocolIds
|
||||
elif isinstance(capability, KernelFactory):
|
||||
elif isinstance(capability, CoreRuntimeFactory):
|
||||
configurationTypes = tuple(capability.configurationTypes)
|
||||
kernelTypes = tuple(capability.kernelTypes)
|
||||
runtimeTypes = _runtimeTypes(capability)
|
||||
|
||||
if not configurationTypes:
|
||||
raise ValueError(
|
||||
f'kernel factory {capability.factoryId!r} must declare '
|
||||
f'core runtime factory {capability.factoryId!r} must declare '
|
||||
f'configuration types'
|
||||
)
|
||||
|
||||
@@ -262,16 +267,16 @@ class PluginRegistry:
|
||||
localConfigurationTypes,
|
||||
),
|
||||
(
|
||||
kernelTypes,
|
||||
'kernel',
|
||||
tuple(self._kernelFactories),
|
||||
localKernelTypes,
|
||||
runtimeTypes,
|
||||
'runtime',
|
||||
tuple(self._runtimeFactories),
|
||||
localRuntimeTypes,
|
||||
),
|
||||
):
|
||||
for itemType in values:
|
||||
if not isinstance(itemType, type):
|
||||
raise TypeError(
|
||||
f'kernel factory {label} types must be classes'
|
||||
f'core runtime factory {label} types must be classes'
|
||||
)
|
||||
|
||||
if any(
|
||||
@@ -286,38 +291,38 @@ class PluginRegistry:
|
||||
|
||||
local.append(itemType)
|
||||
|
||||
detail = (configurationTypes, kernelTypes)
|
||||
detail = (configurationTypes, runtimeTypes)
|
||||
elif isinstance(capability, TrafficStatsProvider):
|
||||
kernelTypes = tuple(capability.kernelTypes)
|
||||
runtimeTypes = _runtimeTypes(capability)
|
||||
|
||||
if not kernelTypes:
|
||||
if not runtimeTypes:
|
||||
raise ValueError(
|
||||
f'traffic stats provider {capability.providerId!r} must '
|
||||
f'declare kernel types'
|
||||
f'declare runtime types'
|
||||
)
|
||||
|
||||
for kernelType in kernelTypes:
|
||||
if not isinstance(kernelType, type):
|
||||
for runtimeType in runtimeTypes:
|
||||
if not isinstance(runtimeType, type):
|
||||
raise TypeError(
|
||||
'traffic stats provider kernel types must be classes'
|
||||
'traffic stats provider runtime types must be classes'
|
||||
)
|
||||
|
||||
if any(
|
||||
issubclass(kernelType, registeredType)
|
||||
or issubclass(registeredType, kernelType)
|
||||
issubclass(runtimeType, registeredType)
|
||||
or issubclass(registeredType, runtimeType)
|
||||
for registeredType in (
|
||||
*self._trafficStatsProviders,
|
||||
*localTrafficStatsKernelTypes,
|
||||
*localTrafficStatsRuntimeTypes,
|
||||
)
|
||||
):
|
||||
raise ValueError(
|
||||
f'traffic stats kernel type '
|
||||
f'{kernelType.__name__!r} overlaps a registered type'
|
||||
f'traffic stats runtime type '
|
||||
f'{runtimeType.__name__!r} overlaps a registered type'
|
||||
)
|
||||
|
||||
localTrafficStatsKernelTypes.append(kernelType)
|
||||
localTrafficStatsRuntimeTypes.append(runtimeType)
|
||||
|
||||
detail = kernelTypes
|
||||
detail = runtimeTypes
|
||||
elif isinstance(capability, SubscriptionDecoder):
|
||||
if not isinstance(capability.priority, int):
|
||||
raise TypeError('subscription decoder priority must be an integer')
|
||||
@@ -356,17 +361,17 @@ class PluginRegistry:
|
||||
|
||||
for protocolId in detail:
|
||||
self._protocolEditors[protocolId] = entry
|
||||
elif isinstance(capability, KernelFactory):
|
||||
configurationTypes, kernelTypes = detail
|
||||
elif isinstance(capability, CoreRuntimeFactory):
|
||||
configurationTypes, runtimeTypes = detail
|
||||
self._factories[capabilityId] = entry
|
||||
|
||||
for configType in configurationTypes:
|
||||
self._configurationFactories[configType] = entry
|
||||
for kernelType in kernelTypes:
|
||||
self._kernelFactories[kernelType] = entry
|
||||
for runtimeType in runtimeTypes:
|
||||
self._runtimeFactories[runtimeType] = entry
|
||||
elif isinstance(capability, TrafficStatsProvider):
|
||||
for kernelType in detail:
|
||||
self._trafficStatsProviders[kernelType] = entry
|
||||
for runtimeType in detail:
|
||||
self._trafficStatsProviders[runtimeType] = entry
|
||||
elif isinstance(capability, SubscriptionDecoder):
|
||||
self._decoders[capabilityId] = entry
|
||||
|
||||
@@ -421,7 +426,7 @@ class PluginRegistry:
|
||||
'_protocolEditors',
|
||||
'_factories',
|
||||
'_configurationFactories',
|
||||
'_kernelFactories',
|
||||
'_runtimeFactories',
|
||||
'_trafficStatsProviders',
|
||||
'_decoders',
|
||||
):
|
||||
@@ -500,9 +505,9 @@ class PluginRegistry:
|
||||
"""Return registered protocol editor providers."""
|
||||
return self.capabilities(CapabilityKind.ProtocolEditor)
|
||||
|
||||
def kernelFactories(self):
|
||||
"""Return registered runtime kernel factories."""
|
||||
return self.capabilities(CapabilityKind.KernelFactory)
|
||||
def coreRuntimeFactories(self):
|
||||
"""Return registered core-runtime factories."""
|
||||
return self.capabilities(CapabilityKind.CoreRuntimeFactory)
|
||||
|
||||
def subscriptionDecoders(self):
|
||||
"""Return subscription decoders in auto-detection priority order."""
|
||||
@@ -571,32 +576,32 @@ class PluginRegistry:
|
||||
|
||||
return None
|
||||
|
||||
def factoryForKernel(self, kernel):
|
||||
"""Return the runtime factory that owns *kernel*."""
|
||||
for kernelType, (_plugin, factory) in self._kernelFactories.items():
|
||||
if isinstance(kernel, kernelType):
|
||||
def runtimeFactoryFor(self, runtime):
|
||||
"""Return the factory that owns *runtime*."""
|
||||
for runtimeType, (_plugin, factory) in self._runtimeFactories.items():
|
||||
if isinstance(runtime, runtimeType):
|
||||
return factory
|
||||
|
||||
return None
|
||||
|
||||
def trafficStatsProviderForKernel(self, kernel):
|
||||
"""Return the traffic-statistics provider that owns *kernel*."""
|
||||
for kernelType, (_plugin, provider) in self._trafficStatsProviders.items():
|
||||
if isinstance(kernel, kernelType):
|
||||
def trafficStatsProviderForRuntime(self, runtime):
|
||||
"""Return the traffic-statistics provider that owns *runtime*."""
|
||||
for runtimeType, (_plugin, provider) in self._trafficStatsProviders.items():
|
||||
if isinstance(runtime, runtimeType):
|
||||
return provider
|
||||
|
||||
return None
|
||||
|
||||
def trafficStatsMonitorForKernels(self, kernels):
|
||||
"""Return the first available monitor for the active runtime kernels."""
|
||||
for kernel in kernels:
|
||||
provider = self.trafficStatsProviderForKernel(kernel)
|
||||
def trafficStatsMonitorForRuntimes(self, runtimes):
|
||||
"""Return the first monitor available for the active core runtimes."""
|
||||
for runtime in runtimes:
|
||||
provider = self.trafficStatsProviderForRuntime(runtime)
|
||||
|
||||
if provider is None:
|
||||
continue
|
||||
|
||||
try:
|
||||
monitor = provider.monitorForKernel(kernel)
|
||||
monitor = provider.monitorForRuntime(runtime)
|
||||
except Exception as ex:
|
||||
# Any non-exit exceptions
|
||||
|
||||
@@ -652,10 +657,10 @@ class PluginRegistry:
|
||||
else None
|
||||
)
|
||||
|
||||
def pluginForKernel(self, kernel):
|
||||
"""Return the plugin that contributes a kernel's factory."""
|
||||
for kernelType, (plugin, _factory) in self._kernelFactories.items():
|
||||
if isinstance(kernel, kernelType):
|
||||
def pluginForRuntime(self, runtime):
|
||||
"""Return the plugin that contributes a runtime's factory."""
|
||||
for runtimeType, (plugin, _factory) in self._runtimeFactories.items():
|
||||
if isinstance(runtime, runtimeType):
|
||||
return plugin
|
||||
|
||||
return None
|
||||
@@ -771,7 +776,7 @@ class PluginRegistry:
|
||||
if result is not None:
|
||||
if not isinstance(result, factory.configurationTypes):
|
||||
logger.error(
|
||||
f'kernel factory {factory.factoryId!r} returned an '
|
||||
f'core runtime factory {factory.factoryId!r} returned an '
|
||||
f'unowned mapping result'
|
||||
)
|
||||
continue
|
||||
@@ -780,7 +785,9 @@ class PluginRegistry:
|
||||
|
||||
if len(factoryMatches) > 1:
|
||||
names = ', '.join(repr(item[0].factoryId) for item in factoryMatches)
|
||||
raise ValueError(f'kernel configuration mapping is ambiguous: {names}')
|
||||
raise ValueError(
|
||||
f'core runtime configuration mapping is ambiguous: {names}'
|
||||
)
|
||||
|
||||
return factoryMatches[0][1] if factoryMatches else None
|
||||
|
||||
@@ -901,7 +908,7 @@ class PluginRegistry:
|
||||
handled = factory.prepareTUN(_connectionOf(config))
|
||||
|
||||
if not isinstance(handled, bool):
|
||||
raise TypeError('kernel TUN preparation result must be a boolean')
|
||||
raise TypeError('core runtime TUN preparation result must be a boolean')
|
||||
|
||||
return handled
|
||||
except TUNPreparationError:
|
||||
@@ -926,7 +933,9 @@ class PluginRegistry:
|
||||
enabled = factory.usesApplicationTun2socks(_connectionOf(config))
|
||||
|
||||
if not isinstance(enabled, bool):
|
||||
raise TypeError('kernel application tun2socks result must be a boolean')
|
||||
raise TypeError(
|
||||
'core runtime application tun2socks result must be a boolean'
|
||||
)
|
||||
|
||||
return enabled
|
||||
except Exception as ex:
|
||||
@@ -953,7 +962,7 @@ class PluginRegistry:
|
||||
for option in options:
|
||||
if not isinstance(option, RoutingOption):
|
||||
raise TypeError(
|
||||
'kernel routing options must be RoutingOption values'
|
||||
'core runtime routing options must be RoutingOption values'
|
||||
)
|
||||
|
||||
if not isinstance(option.id, str) or not option.id.strip():
|
||||
@@ -993,14 +1002,14 @@ class PluginRegistry:
|
||||
|
||||
return routing if routing in optionIds else optionIds[0]
|
||||
|
||||
def createKernel(self, config, routing, **kwargs):
|
||||
"""Create a prepared kernel launch for *config*."""
|
||||
def createCoreRuntime(self, config, routing, **kwargs):
|
||||
"""Create a prepared core-runtime launch for *config*."""
|
||||
factory = self.factoryForConfig(config)
|
||||
|
||||
if factory is None:
|
||||
return None
|
||||
|
||||
request = KernelRequest(
|
||||
request = CoreRuntimeRequest(
|
||||
configuration=_connectionOf(config),
|
||||
routing=self.normalizeRouting(config, routing),
|
||||
exitCallback=kwargs.pop('exitCallback', None),
|
||||
@@ -1014,35 +1023,42 @@ class PluginRegistry:
|
||||
if launch is None:
|
||||
return None
|
||||
|
||||
if not isinstance(launch, KernelLaunch):
|
||||
raise TypeError('kernel factory must return a KernelLaunch value')
|
||||
|
||||
if factory.kernelTypes and not isinstance(launch.kernel, factory.kernelTypes):
|
||||
if not isinstance(launch, CoreRuntimeLaunch):
|
||||
raise TypeError(
|
||||
f'kernel factory {factory.factoryId!r} returned an unowned kernel'
|
||||
'core runtime factory must return a CoreRuntimeLaunch value'
|
||||
)
|
||||
|
||||
runtimeTypes = _runtimeTypes(factory)
|
||||
|
||||
if runtimeTypes and not isinstance(launch.runtime, runtimeTypes):
|
||||
raise TypeError(
|
||||
f'core runtime factory {factory.factoryId!r} returned an '
|
||||
f'unowned runtime'
|
||||
)
|
||||
|
||||
return launch
|
||||
|
||||
def startKernel(self, config, routing, **kwargs):
|
||||
"""Create and start the runtime kernel selected for *config*."""
|
||||
def startCoreRuntime(self, config, routing, **kwargs):
|
||||
"""Create and start the core runtime selected for *config*."""
|
||||
try:
|
||||
launch = self.createKernel(config, routing, **kwargs)
|
||||
launch = self.createCoreRuntime(config, routing, **kwargs)
|
||||
|
||||
return (
|
||||
(launch.kernel, launch.start()) if launch is not None else (None, False)
|
||||
(launch.runtime, launch.start())
|
||||
if launch is not None
|
||||
else (None, False)
|
||||
)
|
||||
except Exception as ex:
|
||||
# Any non-exit exceptions
|
||||
|
||||
factory = self.factoryForConfig(config)
|
||||
factoryId = factory.factoryId if factory is not None else 'unknown'
|
||||
logger.error(f'kernel start failed for {factoryId!r}: {ex}')
|
||||
logger.error(f'core runtime start failed for {factoryId!r}: {ex}')
|
||||
|
||||
return None, False
|
||||
|
||||
def prepareDownloadTest(self, config, port: int):
|
||||
"""Create a proxy-only test configuration through its kernel factory."""
|
||||
"""Create a proxy-only test configuration through its runtime factory."""
|
||||
factory = self.factoryForConfig(config)
|
||||
|
||||
return (
|
||||
@@ -1101,8 +1117,8 @@ class PluginRegistry:
|
||||
return None
|
||||
|
||||
def configureEnvironment(self):
|
||||
"""Allow every kernel factory to configure its process environment."""
|
||||
for factory in self.kernelFactories():
|
||||
"""Allow every runtime factory to configure its execution environment."""
|
||||
for factory in self.coreRuntimeFactories():
|
||||
try:
|
||||
factory.configureEnvironment()
|
||||
except Exception as ex:
|
||||
@@ -1111,10 +1127,10 @@ class PluginRegistry:
|
||||
logger.error(f'environment hook failed for {factory.factoryId!r}: {ex}')
|
||||
|
||||
def coreVersions(self):
|
||||
"""Return version strings reported by every kernel factory."""
|
||||
"""Return version strings reported by every core-runtime factory."""
|
||||
versions = []
|
||||
|
||||
for factory in self.kernelFactories():
|
||||
for factory in self.coreRuntimeFactories():
|
||||
try:
|
||||
versions.extend(factory.coreVersions())
|
||||
except Exception as ex:
|
||||
@@ -1127,10 +1143,10 @@ class PluginRegistry:
|
||||
return tuple(filter(None, versions))
|
||||
|
||||
def logTimestampPatterns(self):
|
||||
"""Return timestamp expressions contributed by all kernel factories."""
|
||||
"""Return timestamp expressions contributed by runtime factories."""
|
||||
patterns = []
|
||||
|
||||
for factory in self.kernelFactories():
|
||||
for factory in self.coreRuntimeFactories():
|
||||
try:
|
||||
patterns.extend(factory.logTimestampPatterns())
|
||||
except Exception as ex:
|
||||
@@ -1144,7 +1160,7 @@ class PluginRegistry:
|
||||
|
||||
def coreExitMessage(self, core, exitcode: int):
|
||||
"""Return the owning factory's special exit message, if any."""
|
||||
factory = self.factoryForKernel(core)
|
||||
factory = self.runtimeFactoryFor(core)
|
||||
|
||||
if factory is None:
|
||||
return None
|
||||
@@ -1161,8 +1177,8 @@ class PluginRegistry:
|
||||
return None
|
||||
|
||||
def afterConnected(self, httpProxy=None):
|
||||
"""Notify every kernel factory after a connection succeeds."""
|
||||
for factory in self.kernelFactories():
|
||||
"""Notify every core-runtime factory after a connection succeeds."""
|
||||
for factory in self.coreRuntimeFactories():
|
||||
try:
|
||||
factory.afterConnected(httpProxy)
|
||||
except Exception as ex:
|
||||
|
||||
@@ -23,10 +23,10 @@ from .API import (
|
||||
PLUGIN_API_VERSION,
|
||||
ActionProvider,
|
||||
CapabilityKind,
|
||||
CoreRuntimeFactory,
|
||||
CoreRuntimeLaunch,
|
||||
CoreRuntimeRequest,
|
||||
FuriousPlugin,
|
||||
KernelFactory,
|
||||
KernelLaunch,
|
||||
KernelRequest,
|
||||
NavigationPageDescriptor,
|
||||
NavigationPageProvider,
|
||||
PluginCapability,
|
||||
@@ -71,10 +71,10 @@ __all__ = [
|
||||
'PLUGIN_ENTRY_POINT_GROUP',
|
||||
'ActionProvider',
|
||||
'CapabilityKind',
|
||||
'CoreRuntimeFactory',
|
||||
'CoreRuntimeLaunch',
|
||||
'CoreRuntimeRequest',
|
||||
'FuriousPlugin',
|
||||
'KernelFactory',
|
||||
'KernelLaunch',
|
||||
'KernelRequest',
|
||||
'NavigationPageDescriptor',
|
||||
'NavigationPageProvider',
|
||||
'PluginCapability',
|
||||
|
||||
@@ -77,7 +77,7 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
super().__init__(*args, **kwargs)
|
||||
|
||||
self.uniqueCleanup = False
|
||||
self.processesPool = list()
|
||||
self.runtimes = list()
|
||||
self._lastStartError = ''
|
||||
|
||||
@property
|
||||
@@ -85,7 +85,7 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
"""Return the concise failure reported by the latest runtime start."""
|
||||
return self._lastStartError
|
||||
|
||||
def _startKernel(
|
||||
def _startCoreRuntime(
|
||||
self,
|
||||
config,
|
||||
routing,
|
||||
@@ -94,9 +94,9 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
proxyModeOnly=False,
|
||||
log=True,
|
||||
**kwargs,
|
||||
) -> Tuple[Union[CoreProcess, None], bool]:
|
||||
"""Construct and start the runtime kernel selected for a configuration."""
|
||||
return getPluginRegistry().startKernel(
|
||||
) -> Tuple[Union[CoreRuntime, None], bool]:
|
||||
"""Construct and start the core runtime selected for a configuration."""
|
||||
return getPluginRegistry().startCoreRuntime(
|
||||
config,
|
||||
routing,
|
||||
exitCallback=exitCallback,
|
||||
@@ -178,10 +178,10 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
else:
|
||||
logger.info(
|
||||
'application-managed tun2socks skipped by the active '
|
||||
'kernel configuration'
|
||||
'core runtime configuration'
|
||||
)
|
||||
|
||||
process, success = self._startKernel(
|
||||
runtime, success = self._startCoreRuntime(
|
||||
configcopy,
|
||||
routing,
|
||||
exitCallback,
|
||||
@@ -191,18 +191,18 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
**kwargs,
|
||||
)
|
||||
|
||||
if process is not None:
|
||||
self.processesPool.append(process)
|
||||
if runtime is not None:
|
||||
self.runtimes.append(runtime)
|
||||
|
||||
if not success:
|
||||
if isinstance(process, CoreProcess):
|
||||
startError = getattr(process, 'startError', None)
|
||||
if isinstance(runtime, CoreRuntime):
|
||||
startError = getattr(runtime, 'startError', None)
|
||||
|
||||
if callable(startError):
|
||||
self._lastStartError = startError()
|
||||
|
||||
if isinstance(process, CoreProcessWorker):
|
||||
logger.error(f'core {process.name()} start failed')
|
||||
if isinstance(runtime, CoreProcessWorker):
|
||||
logger.error(f'core {runtime.name()} start failed')
|
||||
|
||||
self.stopAll()
|
||||
|
||||
@@ -260,7 +260,7 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
return abortStart(f'unrecognized platform: {PLATFORM}')
|
||||
|
||||
tun = Tun2socks(exitCallback=exitCallback, msgCallback=msgCallbackTUN_)
|
||||
self.processesPool.append(tun)
|
||||
self.runtimes.append(tun)
|
||||
|
||||
tcpSendBufferSize, tcpReceiveBufferSize, tcpAutoTuning = (
|
||||
userTcpSendBufferSize(),
|
||||
@@ -566,28 +566,28 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
return True
|
||||
|
||||
def allRunning(self) -> bool:
|
||||
"""Return whether every managed process is running."""
|
||||
return all(process.isAlive() for process in self.processesPool)
|
||||
"""Return whether every managed core runtime is running."""
|
||||
return all(runtime.isAlive() for runtime in self.runtimes)
|
||||
|
||||
def anyRunning(self) -> bool:
|
||||
"""Return whether any managed process is running."""
|
||||
return any(process.isAlive() for process in self.processesPool)
|
||||
"""Return whether any managed core runtime is running."""
|
||||
return any(runtime.isAlive() for runtime in self.runtimes)
|
||||
|
||||
def stopAll(self):
|
||||
"""Stop every managed proxy and TUN process."""
|
||||
"""Stop every managed proxy-core and TUN runtime."""
|
||||
try:
|
||||
for process in list(self.processesPool):
|
||||
if not isinstance(process, CoreProcess):
|
||||
for runtime in list(self.runtimes):
|
||||
if not isinstance(runtime, CoreRuntime):
|
||||
continue
|
||||
|
||||
try:
|
||||
process.stop()
|
||||
runtime.stop()
|
||||
except Exception as ex:
|
||||
# Any non-exit exceptions
|
||||
|
||||
logger.error(f'error stopping core process: {ex}')
|
||||
logger.error(f'error stopping core runtime: {ex}')
|
||||
finally:
|
||||
dispose = getattr(process, 'dispose', None)
|
||||
dispose = getattr(runtime, 'dispose', None)
|
||||
|
||||
if callable(dispose):
|
||||
try:
|
||||
@@ -595,9 +595,9 @@ class ConnectionManager(Mixins.CleanupOnExit):
|
||||
except Exception as ex:
|
||||
# Any non-exit exceptions
|
||||
|
||||
logger.error(f'error disposing core process: {ex}')
|
||||
logger.error(f'error disposing core runtime: {ex}')
|
||||
finally:
|
||||
self.processesPool.clear()
|
||||
self.runtimes.clear()
|
||||
|
||||
def cleanup(self):
|
||||
"""Release resources owned by the core manager."""
|
||||
|
||||
@@ -42,7 +42,7 @@ def isCoreActive(coreType) -> bool:
|
||||
controller = AppConnectionController()
|
||||
|
||||
return controller.isConnected() and any(
|
||||
isinstance(process, coreType) for process in controller.processes
|
||||
isinstance(runtime, coreType) for runtime in controller.runtimes
|
||||
)
|
||||
except (AttributeError, RuntimeError):
|
||||
return False
|
||||
|
||||
@@ -247,10 +247,10 @@ class TrafficStatsManager(
|
||||
self._requestSample()
|
||||
|
||||
@staticmethod
|
||||
def _activeProcesses():
|
||||
"""Return the processes owned by the active connection controller."""
|
||||
def _activeRuntimes():
|
||||
"""Return the runtimes owned by the active connection controller."""
|
||||
try:
|
||||
return AppConnectionController().processes
|
||||
return AppConnectionController().runtimes
|
||||
except (AttributeError, RuntimeError):
|
||||
return tuple()
|
||||
|
||||
@@ -369,8 +369,8 @@ class TrafficStatsManager(
|
||||
return
|
||||
|
||||
if self._connected:
|
||||
monitor = getPluginRegistry().trafficStatsMonitorForKernels(
|
||||
self._activeProcesses()
|
||||
monitor = getPluginRegistry().trafficStatsMonitorForRuntimes(
|
||||
self._activeRuntimes()
|
||||
)
|
||||
self._activateMonitor(monitor)
|
||||
|
||||
@@ -504,8 +504,8 @@ class TrafficStatsManager(
|
||||
|
||||
return
|
||||
|
||||
monitor = getPluginRegistry().trafficStatsMonitorForKernels(
|
||||
self._activeProcesses()
|
||||
monitor = getPluginRegistry().trafficStatsMonitorForRuntimes(
|
||||
self._activeRuntimes()
|
||||
)
|
||||
self._activateMonitor(monitor)
|
||||
|
||||
|
||||
@@ -561,13 +561,13 @@ class TestDownloadSpeedWorker(WebGETManager):
|
||||
def coreExitCallback(self, config: ConfigFactory, exitcode: int):
|
||||
"""Handle the core exit callback."""
|
||||
try:
|
||||
if exitcode == CoreProcess.ExitCode.ConfigurationError.value:
|
||||
if exitcode == CoreRuntime.ExitCode.ConfigurationError.value:
|
||||
self.factory.metadata.speed = 'Invalid'
|
||||
self.sync()
|
||||
elif exitcode == CoreProcess.ExitCode.ServerStartFailure.value:
|
||||
elif exitcode == CoreRuntime.ExitCode.ServerStartFailure.value:
|
||||
self.factory.metadata.speed = 'Core start failed'
|
||||
self.sync()
|
||||
elif exitcode == CoreProcess.ExitCode.SystemShuttingDown.value:
|
||||
elif exitcode == CoreRuntime.ExitCode.SystemShuttingDown.value:
|
||||
pass
|
||||
else:
|
||||
self.factory.metadata.speed = f'Core exited {exitcode}'
|
||||
@@ -575,7 +575,7 @@ class TestDownloadSpeedWorker(WebGETManager):
|
||||
finally:
|
||||
self.must()
|
||||
|
||||
def _startKernel(self, config) -> bool:
|
||||
def _startCoreRuntime(self, config) -> bool:
|
||||
"""Prepare and start a download test through its runtime factory."""
|
||||
configcopy = getPluginRegistry().prepareDownloadTest(config, self.port)
|
||||
|
||||
@@ -617,7 +617,7 @@ class TestDownloadSpeedWorker(WebGETManager):
|
||||
self.factory.metadata.speed = 'Invalid'
|
||||
self.sync()
|
||||
else:
|
||||
if not self._startKernel(self.factory) or appIsExiting():
|
||||
if not self._startCoreRuntime(self.factory) or appIsExiting():
|
||||
return
|
||||
|
||||
self.configureHttpProxy(f'127.0.0.1:{self.port}')
|
||||
|
||||
Reference in New Issue
Block a user