mirror of
https://github.com/LorenEteval/Furious.git
synced 2026-10-07 06:18:05 +03:00
Add Xray traffic statistics monitoring
Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
@@ -29,6 +29,7 @@ from .Process import *
|
||||
from .ProtocolEditors import XRAY_PROTOCOL_EDITORS
|
||||
from .Protocols import XRAY_PROTOCOL_HANDLERS
|
||||
from .Routing import *
|
||||
from .Stats import *
|
||||
from .TUN import *
|
||||
|
||||
import os
|
||||
@@ -297,10 +298,21 @@ class XrayKernelFactory(KernelFactory):
|
||||
|
||||
config['routing'] = routingObject
|
||||
|
||||
if proxyModeOnly:
|
||||
statsTarget = None
|
||||
else:
|
||||
try:
|
||||
statsTarget = configureXrayStats(config)
|
||||
except Exception as ex:
|
||||
logger.error(f'failed to configure Xray traffic statistics: {ex}')
|
||||
|
||||
statsTarget = None
|
||||
|
||||
process = XrayCore(
|
||||
exitCallback=request.exitCallback,
|
||||
msgCallback=request.messageCallback,
|
||||
)
|
||||
process.xrayStatsTarget = statsTarget
|
||||
|
||||
return KernelLaunch(process, config, options=request.options)
|
||||
|
||||
@@ -389,5 +401,6 @@ class XrayPlugin(FuriousPlugin):
|
||||
*XRAY_PROTOCOL_HANDLERS,
|
||||
*XRAY_PROTOCOL_EDITORS,
|
||||
XrayKernelFactory(),
|
||||
XrayStatsProvider(),
|
||||
XrayActionProvider(),
|
||||
)
|
||||
|
||||
@@ -0,0 +1,279 @@
|
||||
# Copyright (C) 2024–present Loren Eteval & contributors <loren.eteval@proton.me>
|
||||
#
|
||||
# This file is part of Furious.
|
||||
#
|
||||
# This program is free software: you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation, either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
"""Configure and query Xray outbound traffic statistics."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from Furious.Plugins.API import (
|
||||
TrafficCounters,
|
||||
TrafficStatsMonitor,
|
||||
TrafficStatsProvider,
|
||||
)
|
||||
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass
|
||||
from typing import Optional
|
||||
|
||||
import json
|
||||
import socket
|
||||
|
||||
try:
|
||||
import xray
|
||||
except Exception:
|
||||
xray = None
|
||||
|
||||
from .Process import XrayCore
|
||||
|
||||
__all__ = [
|
||||
'XRAY_STATS_QUERY_TIMEOUT',
|
||||
'XrayStatsProvider',
|
||||
'XrayStatsTarget',
|
||||
'buildXrayStats',
|
||||
'configureXrayStats',
|
||||
'queryXrayOutboundStats',
|
||||
'queryXrayStats',
|
||||
]
|
||||
|
||||
XRAY_STATS_API_HOST = '127.0.0.1'
|
||||
XRAY_STATS_QUERY_TIMEOUT = 1
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class XrayStatsTarget:
|
||||
"""Describe the active local Xray statistics API and outbound tag."""
|
||||
|
||||
apiServer: str
|
||||
outboundTag: str
|
||||
timeout: int = XRAY_STATS_QUERY_TIMEOUT
|
||||
reset: bool = False
|
||||
|
||||
|
||||
def _proxyOutboundTag(config: Mapping) -> Optional[str]:
|
||||
"""Return the tag of the configured Xray proxy outbound."""
|
||||
if not isinstance(config, Mapping):
|
||||
return None
|
||||
|
||||
outbounds = config.get('outbounds', [])
|
||||
|
||||
if not isinstance(outbounds, list):
|
||||
return None
|
||||
|
||||
for outbound in outbounds:
|
||||
if not isinstance(outbound, Mapping):
|
||||
continue
|
||||
|
||||
tag = outbound.get('tag')
|
||||
|
||||
if isinstance(tag, str) and tag.casefold() == 'proxy':
|
||||
return tag
|
||||
|
||||
return None
|
||||
|
||||
|
||||
def _availableXrayStatsApiServer() -> str:
|
||||
"""Select an available loopback endpoint for the next Xray launch."""
|
||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as listener:
|
||||
listener.bind((XRAY_STATS_API_HOST, 0))
|
||||
|
||||
return f'{XRAY_STATS_API_HOST}:{listener.getsockname()[1]}'
|
||||
|
||||
|
||||
def buildXrayStats(
|
||||
apiServer: Optional[str] = None,
|
||||
) -> dict:
|
||||
"""Build the Xray API, stats, and outbound policy configuration."""
|
||||
apiServer = apiServer or _availableXrayStatsApiServer()
|
||||
|
||||
return {
|
||||
'api': {
|
||||
'tag': 'api',
|
||||
'listen': apiServer,
|
||||
'services': ['StatsService'],
|
||||
},
|
||||
'stats': {},
|
||||
'policy': {
|
||||
'system': {
|
||||
'statsOutboundUplink': True,
|
||||
'statsOutboundDownlink': True,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def configureXrayStats(
|
||||
config: dict, apiServer: Optional[str] = None
|
||||
) -> Optional[XrayStatsTarget]:
|
||||
"""Enable statistics in *config* for its active proxy outbound."""
|
||||
outboundTag = _proxyOutboundTag(config)
|
||||
|
||||
if outboundTag is None:
|
||||
return None
|
||||
|
||||
currentApi = config.get('api')
|
||||
|
||||
if not isinstance(currentApi, dict):
|
||||
currentApi = {}
|
||||
config['api'] = currentApi
|
||||
|
||||
configuredServer = currentApi.get('listen')
|
||||
if not isinstance(configuredServer, str) or not configuredServer.strip():
|
||||
configuredServer = None
|
||||
|
||||
apiServer = apiServer or configuredServer or _availableXrayStatsApiServer()
|
||||
generated = buildXrayStats(apiServer)
|
||||
|
||||
currentApi.setdefault('tag', generated['api']['tag'])
|
||||
currentApi['listen'] = apiServer
|
||||
|
||||
services = currentApi.get('services', [])
|
||||
if not isinstance(services, list):
|
||||
services = []
|
||||
|
||||
currentApi['services'] = list(
|
||||
dict.fromkeys(
|
||||
[
|
||||
*(service for service in services if isinstance(service, str)),
|
||||
'StatsService',
|
||||
]
|
||||
)
|
||||
)
|
||||
|
||||
if not isinstance(config.get('stats'), dict):
|
||||
config['stats'] = generated['stats']
|
||||
|
||||
policy = config.get('policy')
|
||||
if not isinstance(policy, dict):
|
||||
policy = {}
|
||||
config['policy'] = policy
|
||||
|
||||
systemPolicy = policy.get('system')
|
||||
if not isinstance(systemPolicy, dict):
|
||||
systemPolicy = {}
|
||||
policy['system'] = systemPolicy
|
||||
|
||||
systemPolicy.update(generated['policy']['system'])
|
||||
|
||||
return XrayStatsTarget(apiServer, outboundTag)
|
||||
|
||||
|
||||
def queryXrayStats(
|
||||
apiServer: str,
|
||||
timeout: int,
|
||||
myPattern: str,
|
||||
reset: bool,
|
||||
) -> Optional[str]:
|
||||
"""Query the Xray binding and return its raw response when available."""
|
||||
if (
|
||||
xray is None
|
||||
or not isinstance(apiServer, str)
|
||||
or not apiServer.strip()
|
||||
or not isinstance(timeout, int)
|
||||
or isinstance(timeout, bool)
|
||||
or timeout <= 0
|
||||
or not isinstance(myPattern, str)
|
||||
or not isinstance(reset, bool)
|
||||
):
|
||||
return None
|
||||
|
||||
try:
|
||||
response = xray.queryStats(apiServer, timeout, myPattern, reset)
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
if isinstance(response, bytes):
|
||||
return response.decode('utf-8', 'replace')
|
||||
|
||||
return response if isinstance(response, str) else None
|
||||
|
||||
|
||||
def queryXrayOutboundStats(target: XrayStatsTarget) -> Optional[TrafficCounters]:
|
||||
"""Return cumulative traffic counters for one configured Xray outbound."""
|
||||
if not isinstance(target, XrayStatsTarget) or not target.outboundTag:
|
||||
return None
|
||||
|
||||
prefix = f'outbound>>>{target.outboundTag}>>>traffic>>>'
|
||||
response = queryXrayStats(
|
||||
target.apiServer,
|
||||
target.timeout,
|
||||
prefix,
|
||||
target.reset,
|
||||
)
|
||||
|
||||
try:
|
||||
payload = json.loads(response)
|
||||
except (AttributeError, TypeError, ValueError):
|
||||
return None
|
||||
|
||||
if not isinstance(payload, Mapping):
|
||||
return None
|
||||
|
||||
statistics = payload.get('stat', [])
|
||||
|
||||
if not isinstance(statistics, list):
|
||||
return None
|
||||
|
||||
if not statistics:
|
||||
return TrafficCounters(uplink=0, downlink=0)
|
||||
|
||||
values = {}
|
||||
|
||||
for statistic in statistics:
|
||||
if not isinstance(statistic, Mapping):
|
||||
continue
|
||||
|
||||
name = statistic.get('name')
|
||||
|
||||
if name not in (f'{prefix}uplink', f'{prefix}downlink'):
|
||||
continue
|
||||
|
||||
try:
|
||||
rawValue = statistic.get('value', 0)
|
||||
|
||||
if isinstance(rawValue, bool):
|
||||
continue
|
||||
|
||||
value = int(rawValue)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
|
||||
if value >= 0:
|
||||
values[name.rsplit('>>>', 1)[-1]] = value
|
||||
|
||||
if not values:
|
||||
return None
|
||||
|
||||
return TrafficCounters(
|
||||
uplink=values.get('uplink', 0),
|
||||
downlink=values.get('downlink', 0),
|
||||
)
|
||||
|
||||
|
||||
class XrayStatsProvider(TrafficStatsProvider):
|
||||
"""Expose active Xray outbound traffic through the plugin capability API."""
|
||||
|
||||
providerId = 'official.xray.stats'
|
||||
kernelTypes = (XrayCore,)
|
||||
|
||||
def monitorForKernel(self, kernel) -> Optional[TrafficStatsMonitor]:
|
||||
"""Return a monitor when the active core has a valid outbound target."""
|
||||
target = getattr(kernel, 'xrayStatsTarget', None)
|
||||
|
||||
if not isinstance(target, XrayStatsTarget):
|
||||
return None
|
||||
|
||||
return TrafficStatsMonitor(queryXrayOutboundStats, target)
|
||||
+38
-1
@@ -21,7 +21,7 @@ from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
from enum import Enum
|
||||
from typing import Any, Mapping, Optional, Tuple
|
||||
from typing import Any, Callable, Mapping, Optional, Tuple
|
||||
|
||||
__all__ = [
|
||||
'PLUGIN_API_VERSION',
|
||||
@@ -42,6 +42,9 @@ __all__ = [
|
||||
'SubscriptionDecoder',
|
||||
'SubscriptionItem',
|
||||
'SubscriptionResult',
|
||||
'TrafficCounters',
|
||||
'TrafficStatsMonitor',
|
||||
'TrafficStatsProvider',
|
||||
]
|
||||
|
||||
PLUGIN_API_VERSION = 3
|
||||
@@ -55,6 +58,7 @@ class CapabilityKind(str, Enum):
|
||||
ProtocolEditor = 'protocol-editor'
|
||||
SubscriptionDecoder = 'subscription-decoder'
|
||||
KernelFactory = 'kernel-factory'
|
||||
TrafficStats = 'traffic-stats'
|
||||
Utility = 'utility'
|
||||
|
||||
|
||||
@@ -247,6 +251,39 @@ class SubscriptionDecoder(PluginCapability):
|
||||
return None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TrafficCounters:
|
||||
"""Store cumulative upload and download byte counters."""
|
||||
|
||||
uplink: int
|
||||
downlink: int
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TrafficStatsMonitor:
|
||||
"""Describe an isolated traffic-statistics query operation."""
|
||||
|
||||
query: Callable[[Any], Optional[TrafficCounters]]
|
||||
target: Any
|
||||
|
||||
|
||||
class TrafficStatsProvider(PluginCapability):
|
||||
"""Provide traffic counters for one or more runtime kernel types."""
|
||||
|
||||
capabilityKind = CapabilityKind.TrafficStats
|
||||
providerId = ''
|
||||
kernelTypes = 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."""
|
||||
return None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KernelRequest:
|
||||
"""Describe one runtime-kernel construction request."""
|
||||
|
||||
@@ -83,6 +83,7 @@ class PluginRegistry:
|
||||
self._factories = {}
|
||||
self._configurationFactories = {}
|
||||
self._kernelFactories = {}
|
||||
self._trafficStatsProviders = {}
|
||||
self._decoders = {}
|
||||
self._initializedPlugins = []
|
||||
self._closed = False
|
||||
@@ -146,6 +147,7 @@ class PluginRegistry:
|
||||
localEditorProtocols = set()
|
||||
localConfigurationTypes = []
|
||||
localKernelTypes = []
|
||||
localTrafficStatsKernelTypes = []
|
||||
entries = []
|
||||
|
||||
for capability in capabilities:
|
||||
@@ -276,6 +278,37 @@ class PluginRegistry:
|
||||
local.append(itemType)
|
||||
|
||||
detail = (configurationTypes, kernelTypes)
|
||||
elif isinstance(capability, TrafficStatsProvider):
|
||||
kernelTypes = tuple(capability.kernelTypes)
|
||||
|
||||
if not kernelTypes:
|
||||
raise ValueError(
|
||||
f'traffic stats provider {capability.providerId!r} must '
|
||||
f'declare kernel types'
|
||||
)
|
||||
|
||||
for kernelType in kernelTypes:
|
||||
if not isinstance(kernelType, type):
|
||||
raise TypeError(
|
||||
'traffic stats provider kernel types must be classes'
|
||||
)
|
||||
|
||||
if any(
|
||||
issubclass(kernelType, registeredType)
|
||||
or issubclass(registeredType, kernelType)
|
||||
for registeredType in (
|
||||
*self._trafficStatsProviders,
|
||||
*localTrafficStatsKernelTypes,
|
||||
)
|
||||
):
|
||||
raise ValueError(
|
||||
f'traffic stats kernel type '
|
||||
f'{kernelType.__name__!r} overlaps a registered type'
|
||||
)
|
||||
|
||||
localTrafficStatsKernelTypes.append(kernelType)
|
||||
|
||||
detail = kernelTypes
|
||||
elif isinstance(capability, SubscriptionDecoder):
|
||||
if not isinstance(capability.priority, int):
|
||||
raise TypeError('subscription decoder priority must be an integer')
|
||||
@@ -322,6 +355,9 @@ class PluginRegistry:
|
||||
self._configurationFactories[configType] = entry
|
||||
for kernelType in kernelTypes:
|
||||
self._kernelFactories[kernelType] = entry
|
||||
elif isinstance(capability, TrafficStatsProvider):
|
||||
for kernelType in detail:
|
||||
self._trafficStatsProviders[kernelType] = entry
|
||||
elif isinstance(capability, SubscriptionDecoder):
|
||||
self._decoders[capabilityId] = entry
|
||||
|
||||
@@ -373,6 +409,7 @@ class PluginRegistry:
|
||||
'_factories',
|
||||
'_configurationFactories',
|
||||
'_kernelFactories',
|
||||
'_trafficStatsProviders',
|
||||
'_decoders',
|
||||
):
|
||||
setattr(
|
||||
@@ -464,6 +501,10 @@ class PluginRegistry:
|
||||
)
|
||||
)
|
||||
|
||||
def trafficStatsProviders(self):
|
||||
"""Return registered runtime traffic-statistics providers."""
|
||||
return self.capabilities(CapabilityKind.TrafficStats)
|
||||
|
||||
def handlerForProtocol(self, protocol):
|
||||
"""Return the handler registered for a protocol identifier."""
|
||||
entry = self._protocols.get(_normalizeIdentifier(protocol))
|
||||
@@ -515,6 +556,55 @@ class PluginRegistry:
|
||||
|
||||
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):
|
||||
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)
|
||||
|
||||
if provider is None:
|
||||
continue
|
||||
|
||||
try:
|
||||
monitor = provider.monitorForKernel(kernel)
|
||||
except Exception as ex:
|
||||
logger.error(
|
||||
f'failed to obtain traffic stats monitor from '
|
||||
f'{provider.providerId!r}: {ex}'
|
||||
)
|
||||
|
||||
continue
|
||||
|
||||
if monitor is None:
|
||||
continue
|
||||
|
||||
if not isinstance(monitor, TrafficStatsMonitor):
|
||||
logger.error(
|
||||
f'traffic stats provider {provider.providerId!r} returned '
|
||||
f'an invalid monitor'
|
||||
)
|
||||
|
||||
continue
|
||||
|
||||
if not callable(monitor.query):
|
||||
logger.error(
|
||||
f'traffic stats provider {provider.providerId!r} returned '
|
||||
f'a monitor without a query callable'
|
||||
)
|
||||
|
||||
continue
|
||||
|
||||
return monitor
|
||||
|
||||
return None
|
||||
|
||||
def pluginForProtocol(self, protocol):
|
||||
"""Return the plugin that contributes *protocol*."""
|
||||
entry = self._protocols.get(_normalizeIdentifier(protocol))
|
||||
|
||||
@@ -38,6 +38,9 @@ from .API import (
|
||||
SubscriptionDecoder,
|
||||
SubscriptionItem,
|
||||
SubscriptionResult,
|
||||
TrafficCounters,
|
||||
TrafficStatsMonitor,
|
||||
TrafficStatsProvider,
|
||||
)
|
||||
from .Profile import (
|
||||
blankConfiguration,
|
||||
@@ -77,6 +80,9 @@ __all__ = [
|
||||
'SubscriptionDecoder',
|
||||
'SubscriptionItem',
|
||||
'SubscriptionResult',
|
||||
'TrafficCounters',
|
||||
'TrafficStatsMonitor',
|
||||
'TrafficStatsProvider',
|
||||
'blankConfiguration',
|
||||
'blankProfile',
|
||||
'configurationFromAny',
|
||||
|
||||
@@ -967,6 +967,22 @@ class AppStyleSheet:
|
||||
color: {palette['danger']};
|
||||
}}
|
||||
|
||||
QWidget#TrafficSpeedBadge {{
|
||||
min-height: 24px;
|
||||
padding: 0 2px;
|
||||
border: 1px solid {palette['border']};
|
||||
border-radius: 6px;
|
||||
background-color: {palette['raised']};
|
||||
color: {palette['text_strong']};
|
||||
}}
|
||||
|
||||
QWidget#TrafficSpeedBadge QLabel {{
|
||||
border: none;
|
||||
padding: 0;
|
||||
background-color: transparent;
|
||||
color: {palette['text_strong']};
|
||||
}}
|
||||
|
||||
QSplitter::handle {{
|
||||
background-color: {palette['border']};
|
||||
}}
|
||||
|
||||
@@ -0,0 +1,306 @@
|
||||
# Copyright (C) 2024–present Loren Eteval & contributors <loren.eteval@proton.me>
|
||||
#
|
||||
# This file is part of Furious.
|
||||
#
|
||||
# This program is free software: you can redistribute it and/or modify
|
||||
# it under the terms of the GNU General Public License as published by
|
||||
# the Free Software Foundation, either version 3 of the License, or
|
||||
# (at your option) any later version.
|
||||
#
|
||||
# This program is distributed in the hope that it will be useful,
|
||||
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||
# GNU General Public License for more details.
|
||||
#
|
||||
# You should have received a copy of the GNU General Public License
|
||||
# along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
"""Sample plugin-provided traffic counters without blocking the Qt UI."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from Furious.Frozenlib import APP, Mixins
|
||||
from Furious.Plugins import TrafficCounters, getPluginRegistry
|
||||
|
||||
from PySide6 import QtCore
|
||||
|
||||
import logging
|
||||
import multiprocessing
|
||||
import queue
|
||||
import time
|
||||
|
||||
__all__ = ['TrafficStatsManager', 'formatTrafficSpeed']
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
KIBIBYTE = 1024
|
||||
MEBIBYTE = KIBIBYTE * KIBIBYTE
|
||||
TRAFFIC_STATS_SAMPLE_INTERVAL = 2000
|
||||
TRAFFIC_STATS_RESULT_INTERVAL = 50
|
||||
|
||||
|
||||
def formatTrafficSpeed(bytesPerSecond: float) -> str:
|
||||
"""Format a byte rate using compact binary units."""
|
||||
try:
|
||||
value = max(float(bytesPerSecond), 0.0)
|
||||
except (TypeError, ValueError):
|
||||
value = 0.0
|
||||
|
||||
if value < KIBIBYTE:
|
||||
return '< 1 KiB/s'
|
||||
|
||||
if value >= MEBIBYTE:
|
||||
amount = value / MEBIBYTE
|
||||
unit = 'MiB/s'
|
||||
else:
|
||||
amount = value / KIBIBYTE
|
||||
unit = 'KiB/s'
|
||||
|
||||
formatted = f'{amount:.2f}'.rstrip('0').rstrip('.')
|
||||
|
||||
return f'{formatted} {unit}'
|
||||
|
||||
|
||||
def _queryTrafficStats(requestQueue, resultQueue):
|
||||
"""Run potentially GIL-blocking statistics calls in an isolated process."""
|
||||
while True:
|
||||
request = requestQueue.get()
|
||||
|
||||
if request is None:
|
||||
return
|
||||
|
||||
generation, monitor = request
|
||||
|
||||
try:
|
||||
counters = monitor.query(monitor.target)
|
||||
except Exception:
|
||||
counters = None
|
||||
|
||||
try:
|
||||
resultQueue.put((generation, counters, time.monotonic()))
|
||||
except Exception:
|
||||
return
|
||||
|
||||
|
||||
class TrafficStatsManager(
|
||||
Mixins.ConnectionAware,
|
||||
Mixins.CleanupOnExit,
|
||||
QtCore.QObject,
|
||||
):
|
||||
"""Periodically convert cumulative plugin counters into traffic speeds."""
|
||||
|
||||
speedChanged = QtCore.Signal(float, float)
|
||||
statisticsUnavailable = QtCore.Signal()
|
||||
|
||||
def __init__(self, parent=None):
|
||||
"""Initialize timers and lazy worker-process state."""
|
||||
super().__init__(parent)
|
||||
|
||||
self._context = multiprocessing.get_context()
|
||||
self._requestQueue = None
|
||||
self._resultQueue = None
|
||||
self._workerProcess = None
|
||||
self._monitor = None
|
||||
self._generation = 0
|
||||
self._queryInFlight = False
|
||||
self._previousCounters = None
|
||||
self._previousSampleTime = None
|
||||
|
||||
self._sampleTimer = QtCore.QTimer(self)
|
||||
self._sampleTimer.setInterval(TRAFFIC_STATS_SAMPLE_INTERVAL)
|
||||
self._sampleTimer.timeout.connect(self._requestSample)
|
||||
|
||||
self._resultTimer = QtCore.QTimer(self)
|
||||
self._resultTimer.setInterval(TRAFFIC_STATS_RESULT_INTERVAL)
|
||||
self._resultTimer.timeout.connect(self._consumeResults)
|
||||
|
||||
@staticmethod
|
||||
def _activeProcesses():
|
||||
"""Return the processes owned by the active tray connection."""
|
||||
try:
|
||||
return tuple(APP().systemTray.ConnectAction.coreManager.processesPool)
|
||||
except (AttributeError, RuntimeError):
|
||||
return tuple()
|
||||
|
||||
def _closeWorker(self):
|
||||
"""Stop and release the isolated query worker."""
|
||||
process = self._workerProcess
|
||||
|
||||
if process is not None and process.is_alive():
|
||||
try:
|
||||
self._requestQueue.put_nowait(None)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
process.join(0.2)
|
||||
|
||||
if process.is_alive():
|
||||
process.terminate()
|
||||
process.join(1.0)
|
||||
|
||||
for workerQueue in (self._requestQueue, self._resultQueue):
|
||||
if workerQueue is None:
|
||||
continue
|
||||
|
||||
try:
|
||||
workerQueue.close()
|
||||
workerQueue.cancel_join_thread()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
self._requestQueue = None
|
||||
self._resultQueue = None
|
||||
self._workerProcess = None
|
||||
|
||||
def _ensureWorker(self) -> bool:
|
||||
"""Start the isolated query worker when it is not already alive."""
|
||||
if self._workerProcess is not None and self._workerProcess.is_alive():
|
||||
return True
|
||||
|
||||
self._closeWorker()
|
||||
self._requestQueue = self._context.Queue(maxsize=1)
|
||||
self._resultQueue = self._context.Queue()
|
||||
self._workerProcess = self._context.Process(
|
||||
target=_queryTrafficStats,
|
||||
args=(self._requestQueue, self._resultQueue),
|
||||
daemon=True,
|
||||
)
|
||||
|
||||
try:
|
||||
self._workerProcess.start()
|
||||
except Exception as ex:
|
||||
logger.error(f'failed to start traffic statistics worker: {ex}')
|
||||
self._closeWorker()
|
||||
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
def _resetSamples(self):
|
||||
"""Discard the previous counter baseline and in-flight marker."""
|
||||
self._queryInFlight = False
|
||||
self._previousCounters = None
|
||||
self._previousSampleTime = None
|
||||
|
||||
def _activateMonitor(self, monitor):
|
||||
"""Begin sampling one plugin-provided monitor."""
|
||||
self._generation += 1
|
||||
self._monitor = monitor
|
||||
self._resetSamples()
|
||||
|
||||
if monitor is None:
|
||||
self._closeWorker()
|
||||
self.statisticsUnavailable.emit()
|
||||
|
||||
return
|
||||
|
||||
if not self._ensureWorker():
|
||||
self.statisticsUnavailable.emit()
|
||||
|
||||
return
|
||||
|
||||
self._resultTimer.start()
|
||||
self._sampleTimer.start()
|
||||
self._requestSample()
|
||||
|
||||
@QtCore.Slot()
|
||||
def _requestSample(self):
|
||||
"""Queue one statistics query unless another query is still running."""
|
||||
if self._monitor is None or self._queryInFlight:
|
||||
return
|
||||
|
||||
if not self._ensureWorker():
|
||||
self.statisticsUnavailable.emit()
|
||||
|
||||
return
|
||||
|
||||
try:
|
||||
self._requestQueue.put_nowait((self._generation, self._monitor))
|
||||
except queue.Full:
|
||||
return
|
||||
except Exception as ex:
|
||||
logger.error(f'failed to queue traffic statistics query: {ex}')
|
||||
|
||||
return
|
||||
|
||||
self._queryInFlight = True
|
||||
|
||||
def _updateSpeeds(self, counters: TrafficCounters, sampledAt: float):
|
||||
"""Calculate and emit rates from one cumulative counter sample."""
|
||||
previousCounters = self._previousCounters
|
||||
previousSampleTime = self._previousSampleTime
|
||||
self._previousCounters = counters
|
||||
self._previousSampleTime = sampledAt
|
||||
|
||||
if previousCounters is None or previousSampleTime is None:
|
||||
self.speedChanged.emit(0.0, 0.0)
|
||||
|
||||
return
|
||||
|
||||
elapsed = sampledAt - previousSampleTime
|
||||
|
||||
if elapsed <= 0:
|
||||
self.speedChanged.emit(0.0, 0.0)
|
||||
|
||||
return
|
||||
|
||||
uplinkDelta = max(counters.uplink - previousCounters.uplink, 0)
|
||||
downlinkDelta = max(counters.downlink - previousCounters.downlink, 0)
|
||||
|
||||
self.speedChanged.emit(uplinkDelta / elapsed, downlinkDelta / elapsed)
|
||||
|
||||
@QtCore.Slot()
|
||||
def _consumeResults(self):
|
||||
"""Drain worker results and ignore samples from older connections."""
|
||||
resultQueue = self._resultQueue
|
||||
|
||||
if resultQueue is None:
|
||||
return
|
||||
|
||||
while True:
|
||||
try:
|
||||
generation, counters, sampledAt = resultQueue.get_nowait()
|
||||
except queue.Empty:
|
||||
return
|
||||
except Exception as ex:
|
||||
logger.error(f'failed to read traffic statistics result: {ex}')
|
||||
|
||||
return
|
||||
|
||||
if generation != self._generation:
|
||||
continue
|
||||
|
||||
self._queryInFlight = False
|
||||
|
||||
if not isinstance(counters, TrafficCounters):
|
||||
self._previousCounters = None
|
||||
self._previousSampleTime = None
|
||||
self.statisticsUnavailable.emit()
|
||||
|
||||
continue
|
||||
|
||||
self._updateSpeeds(counters, sampledAt)
|
||||
|
||||
if self._resultQueue is not resultQueue:
|
||||
return
|
||||
|
||||
def connectedCallback(self):
|
||||
"""Discover and activate statistics for the connected runtime."""
|
||||
monitor = getPluginRegistry().trafficStatsMonitorForKernels(
|
||||
self._activeProcesses()
|
||||
)
|
||||
self._activateMonitor(monitor)
|
||||
|
||||
def disconnectedCallback(self):
|
||||
"""Stop sampling and clear traffic speeds after disconnecting."""
|
||||
self._sampleTimer.stop()
|
||||
self._resultTimer.stop()
|
||||
self._generation += 1
|
||||
self._monitor = None
|
||||
self._resetSamples()
|
||||
self._closeWorker()
|
||||
self.statisticsUnavailable.emit()
|
||||
|
||||
def cleanup(self):
|
||||
"""Release timers and the isolated query worker."""
|
||||
self.disconnectedCallback()
|
||||
@@ -27,6 +27,7 @@ from .SubscriptionImporter import (
|
||||
SubscriptionImportService,
|
||||
SubscriptionSource,
|
||||
)
|
||||
from .TrafficStatsManager import TrafficStatsManager, formatTrafficSpeed
|
||||
from .UpdateManager import UpdateManager
|
||||
|
||||
__all__ = [
|
||||
@@ -36,5 +37,7 @@ __all__ = [
|
||||
'SubscriptionImportResult',
|
||||
'SubscriptionImportService',
|
||||
'SubscriptionSource',
|
||||
'TrafficStatsManager',
|
||||
'UpdateManager',
|
||||
'formatTrafficSpeed',
|
||||
]
|
||||
|
||||
@@ -26,7 +26,12 @@ from Furious.Repository import *
|
||||
from Furious.Plugins import CapabilityKind, getPluginRegistry
|
||||
from Furious.Qt import *
|
||||
from Furious.Qt import gettext as _
|
||||
from Furious.Service import ConnectivityManager, UpdateManager
|
||||
from Furious.Service import (
|
||||
ConnectivityManager,
|
||||
TrafficStatsManager,
|
||||
UpdateManager,
|
||||
formatTrafficSpeed,
|
||||
)
|
||||
from Furious.Actions.Import import (
|
||||
ImportFromFileAction,
|
||||
ImportJSONFromClipboardAction,
|
||||
@@ -253,6 +258,85 @@ class NetworkStateBadge(Mixins.QTranslatable, Mixins.ThemeAware, QWidget):
|
||||
self.updateStatusText()
|
||||
|
||||
|
||||
class TrafficSpeedBadge(Mixins.ThemeAware, QWidget):
|
||||
"""Display independently updated upload and download speeds."""
|
||||
|
||||
UploadIconFileName = 'cloud-arrow-up.svg'
|
||||
DownloadIconFileName = 'cloud-arrow-down.svg'
|
||||
IconSize = QtCore.QSize(16, 16)
|
||||
|
||||
def __init__(self, parent=None):
|
||||
"""Initialize dedicated traffic-direction icons and labels."""
|
||||
super().__init__(parent)
|
||||
|
||||
self.setObjectName('TrafficSpeedBadge')
|
||||
self.setVisible(False)
|
||||
|
||||
self.uploadIconLabel = QLabel(parent=self)
|
||||
self.uploadTextLabel = AppQLabel(translatable=False, parent=self)
|
||||
self.downloadIconLabel = QLabel(parent=self)
|
||||
self.downloadTextLabel = AppQLabel(translatable=False, parent=self)
|
||||
|
||||
self._layout = QHBoxLayout(self)
|
||||
self._layout.setContentsMargins(8, 3, 8, 3)
|
||||
self._layout.setSpacing(6)
|
||||
self._layout.addWidget(self.uploadIconLabel)
|
||||
self._layout.addWidget(self.uploadTextLabel)
|
||||
self._layout.addSpacing(4)
|
||||
self._layout.addWidget(self.downloadIconLabel)
|
||||
self._layout.addWidget(self.downloadTextLabel)
|
||||
|
||||
self.setIconByTheme(APP().theme())
|
||||
|
||||
def setIconByTheme(self, theme: str):
|
||||
"""Apply theme-appropriate upload and download icons."""
|
||||
iconFactory = (
|
||||
bootstrapIconWhite if theme == AppStyleSheet.Dark else bootstrapIcon
|
||||
)
|
||||
self.uploadIconLabel.setPixmap(
|
||||
iconFactory(self.UploadIconFileName).pixmap(self.IconSize)
|
||||
)
|
||||
self.downloadIconLabel.setPixmap(
|
||||
iconFactory(self.DownloadIconFileName).pixmap(self.IconSize)
|
||||
)
|
||||
|
||||
@QtCore.Slot(float, float)
|
||||
def setSpeeds(self, upload: float, download: float):
|
||||
"""Display formatted upload and download byte rates."""
|
||||
self.uploadTextLabel.setText(formatTrafficSpeed(upload))
|
||||
self.downloadTextLabel.setText(formatTrafficSpeed(download))
|
||||
self.setVisible(True)
|
||||
|
||||
@QtCore.Slot()
|
||||
def clearSpeeds(self):
|
||||
"""Hide stale speeds when statistics are unavailable."""
|
||||
self.uploadTextLabel.clear()
|
||||
self.downloadTextLabel.clear()
|
||||
self.setVisible(False)
|
||||
|
||||
def themeChangedCallback(self, theme: str):
|
||||
"""Refresh traffic-direction icons after a theme change."""
|
||||
self.setIconByTheme(theme)
|
||||
|
||||
|
||||
class ConnectionStatusWidget(QWidget):
|
||||
"""Compose network state and traffic speed as one permanent status item."""
|
||||
|
||||
def __init__(self, parent=None):
|
||||
"""Initialize separate network-state and traffic-speed components."""
|
||||
super().__init__(parent)
|
||||
|
||||
self.setObjectName('ConnectionStatusWidget')
|
||||
self.networkState = NetworkStateBadge(parent=self)
|
||||
self.trafficSpeed = TrafficSpeedBadge(parent=self)
|
||||
|
||||
self._layout = QHBoxLayout(self)
|
||||
self._layout.setContentsMargins(0, 0, 0, 0)
|
||||
self._layout.setSpacing(6)
|
||||
self._layout.addWidget(self.networkState)
|
||||
self._layout.addWidget(self.trafficSpeed)
|
||||
|
||||
|
||||
class SearchButton(AppQPushButton):
|
||||
"""Represent search button."""
|
||||
|
||||
@@ -726,9 +810,16 @@ class MainWindow(AppQMainWindow):
|
||||
# TODO: Custom status tip
|
||||
# self.setStatusBar(QStatusBar(self))
|
||||
|
||||
self.networkState = NetworkStateBadge(parent=self)
|
||||
self.connectionStatus = ConnectionStatusWidget(parent=self)
|
||||
self.networkState = self.connectionStatus.networkState
|
||||
self.trafficSpeed = self.connectionStatus.trafficSpeed
|
||||
self.trafficStatsManager = TrafficStatsManager(parent=self)
|
||||
self.trafficStatsManager.speedChanged.connect(self.trafficSpeed.setSpeeds)
|
||||
self.trafficStatsManager.statisticsUnavailable.connect(
|
||||
self.trafficSpeed.clearSpeeds
|
||||
)
|
||||
|
||||
self.statusBar().addPermanentWidget(self.networkState)
|
||||
self.statusBar().addPermanentWidget(self.connectionStatus)
|
||||
|
||||
self._widget = QWidget()
|
||||
self._layout = QVBoxLayout(self._widget)
|
||||
|
||||
Reference in New Issue
Block a user