refactor: extract profile testing service

Move latency and download-speed scheduling out of ServerTableView into an identity-safe service. Add bounded Ping execution, a dedicated adaptive Qt TCPing thread, concurrent non-blocking core startup, and subscription-scoped cancellation with stale-result rejection. Cover scheduler behavior, cancellation races, repository reconciliation, and Qt resource cleanup.

Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
Loren Eteval
2026-08-30 12:46:43 +08:00
parent 2ead3fe06e
commit 3844a8b9aa
11 changed files with 2783 additions and 698 deletions
+5 -3
View File
@@ -27,7 +27,7 @@ from PySide6 import QtCore
from abc import ABC
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Dict, Tuple, Union
from typing import Any, Callable, ClassVar, Dict, Tuple, Union
import os
import sys
@@ -66,12 +66,14 @@ class CoreProcessState(Enum):
class CoreLaunchSpec:
"""Describe the parameters required by a core launch operation."""
DefaultWaitTime: ClassVar[int] = 2500
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
waitTime: int = DefaultWaitTime
def __post_init__(self):
"""Normalize the initialized core launch spec values."""
@@ -90,7 +92,7 @@ class CoreLaunchSpec:
daemon, waitCore, waitTime, target, args = (
kwargs.pop('daemon', True),
kwargs.pop('waitCore', True),
kwargs.pop('waitTime', 2500),
kwargs.pop('waitTime', cls.DefaultWaitTime),
kwargs.pop('target', None),
kwargs.pop('args', tuple()),
)
+14
View File
@@ -33,6 +33,20 @@
repository or UI mutation into a decoder merely to shorten this path.
- Endpoint inspection uses only the active proxy, neutral request metadata, bounded caches, and connection generations.
Reject stale results and disclose actual providers without exposing profile credentials or complete destinations.
- Tcping execution owns sockets and deadlines in one lazily started Qt networking thread. Accept only endpoint-level
requests, keep the adaptive concurrency window bounded, return results through thread-safe events, and synchronously
destroy the engine in its own thread after cancellation and event-loop shutdown.
- `ProfileTestManager` is the sole profile-test write-back boundary. Jobs combine a stable ID/fingerprint/snapshot with
explicit latency or download options; workers return results without mutating profiles, and the manager resolves the
current identity before changing metadata or notifying presentation. Repository mutation refreshes its identity map;
subscription commit invalidates only that group's work and clears only that group's current results.
- Keep Ping, Tcping, and download execution specialized beneath that owner: blocking Ping uses a private bounded pool
and event-based completion; Tcping deduplicates endpoints and fans out results in bounded GUI batches; serial and
concurrent downloads share one scheduler implementation with separate concurrency/port ranges. Concurrent download
admission launches temporary cores without a synchronous readiness wait; each worker owns the readiness timer before
starting HTTP. Cancellation during core launch, readiness, or network startup must defer terminal publication and
deletion until an active startup frame unwinds, and shutdown must stop each exact pool, thread, reply, timer, and
temporary runtime.
## Verification
File diff suppressed because it is too large Load Diff
+6
View File
@@ -103,6 +103,7 @@ class SubscriptionManager(HttpGetManager):
"""Own subscription networking, decoding, reconciliation, and persistence."""
subscriptionsChanged = QtCore.Signal()
subscriptionCommitted = QtCore.Signal(str)
updateCompleted = QtCore.Signal(object)
def __init__(self, parent=None, **kwargs):
@@ -456,6 +457,11 @@ class SubscriptionManager(HttpGetManager):
committedSuccess.append(param)
# Reconciliation is the commit boundary. Notify consumers now so
# subscription-scoped work is invalidated before the rest of a
# multi-subscription batch finishes.
self.subscriptionCommitted.emit(param['unique'])
self._recordGroupSuccess(param, result)
for param in failureArgs:
+347
View File
@@ -0,0 +1,347 @@
# 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/>.
"""Provide bounded asynchronous TCP latency probes in one Qt network thread."""
from __future__ import annotations
from Furious.Qt.Signals import connectWeakly
from PySide6 import QtCore
from PySide6.QtNetwork import QAbstractSocket, QTcpSocket
from shiboken6 import delete as deleteQObject
from dataclasses import dataclass
import collections
import weakref
__all__ = [
'TcpingCancelEvent',
'TcpingEngine',
'TcpingRequest',
'TcpingRequestBatchEvent',
'TcpingResultEvent',
'TcpingThread',
]
TCPING_REQUEST_EVENT_TYPE, TCPING_CANCEL_EVENT_TYPE, TCPING_RESULT_EVENT_TYPE = (
QtCore.QEvent.Type(QtCore.QEvent.registerEventType()),
QtCore.QEvent.Type(QtCore.QEvent.registerEventType()),
QtCore.QEvent.Type(QtCore.QEvent.registerEventType()),
)
@dataclass(frozen=True)
class TcpingRequest:
"""Describe one unique endpoint probe without retaining a live profile."""
requestId: int
address: str
port: int
timeoutMilliseconds: int = 2000
class TcpingRequestBatchEvent(QtCore.QEvent):
"""Carry one coalesced request batch into the networking thread."""
def __init__(self, requests):
"""Initialize the TcpingRequestBatchEvent."""
super().__init__(TCPING_REQUEST_EVENT_TYPE)
self.requests = tuple(requests)
class TcpingCancelEvent(QtCore.QEvent):
"""Carry request cancellation into the networking thread."""
def __init__(self, requestIds):
"""Initialize the TcpingCancelEvent."""
super().__init__(TCPING_CANCEL_EVENT_TYPE)
self.requestIds = frozenset(requestIds)
class TcpingResultEvent(QtCore.QEvent):
"""Carry one probe result back to its GUI-thread scheduler."""
def __init__(self, requestId: int, result):
"""Initialize the TcpingResultEvent."""
super().__init__(TCPING_RESULT_EVENT_TYPE)
self.requestId = requestId
self.result = result
class TcpingProbe(QtCore.QObject):
"""Own one asynchronous socket and timeout inside the networking thread."""
finished = QtCore.Signal(int, object, bool, bool)
def __init__(self, request: TcpingRequest, parent=None):
"""Initialize the TcpingProbe."""
super().__init__(parent)
self.request = request
self.completionHasRun = False
self.elapsedTimer = QtCore.QElapsedTimer()
self.socket = QTcpSocket(self)
self.timeoutTimer = QtCore.QTimer(self)
self.timeoutTimer.setSingleShot(True)
connectWeakly(
self.socket.connected,
self,
'handleConnected',
sender=self.socket,
)
connectWeakly(
self.socket.errorOccurred,
self,
'handleSocketError',
sender=self.socket,
)
connectWeakly(
self.timeoutTimer.timeout,
self,
'handleTimeout',
sender=self.timeoutTimer,
)
def start(self):
"""Start one non-blocking TCP connection attempt."""
self.elapsedTimer.start()
self.timeoutTimer.start(self.request.timeoutMilliseconds)
self.socket.connectToHost(self.request.address, self.request.port)
@QtCore.Slot()
def handleConnected(self):
"""Publish elapsed time after Qt completes the TCP handshake."""
milliseconds = round(self.elapsedTimer.nsecsElapsed() / 1_000_000)
self.complete(f'{max(milliseconds, 0)}ms')
@QtCore.Slot(QAbstractSocket.SocketError)
def handleSocketError(self, _error):
"""Treat a rejected or unreachable endpoint as a timed-out probe."""
self.complete('Timeout')
@QtCore.Slot()
def handleTimeout(self):
"""Abort a connection attempt that exceeded the bounded deadline."""
self.complete('Timeout', deadlineExpired=True)
def complete(self, result, *, deadlineExpired=False, cancelled=False):
"""Stop owned resources and publish exactly one terminal outcome."""
if self.completionHasRun:
return
self.completionHasRun = True
self.timeoutTimer.stop()
if self.socket.state() != QAbstractSocket.SocketState.UnconnectedState:
self.socket.abort()
self.finished.emit(
self.request.requestId,
result,
bool(deadlineExpired),
bool(cancelled),
)
def cancel(self):
"""Abort this probe without publishing a user-visible result."""
self.complete(None, cancelled=True)
class TcpingEngine(QtCore.QObject):
"""Own a bounded adaptive set of probes in one dedicated event loop."""
InitialConcurrency = 2
def __init__(
self,
resultReceiver,
maxConcurrency: int,
parent=None,
**kwargs,
):
"""Initialize the TcpingEngine without constructing thread-bound children."""
super().__init__(parent)
self.resultReceiver = weakref.ref(resultReceiver)
self.maxConcurrency = max(int(maxConcurrency), 1)
self.currentConcurrency = min(self.InitialConcurrency, self.maxConcurrency)
self.probeFactory = kwargs.pop('probeFactory', TcpingProbe)
self.pendingRequests = collections.deque()
self.activeProbes = {}
self.responsiveCompletionCount = 0
self.stopping = False
def event(self, event):
"""Consume request/cancellation commands in this object's thread."""
if event.type() == TCPING_REQUEST_EVENT_TYPE:
self.enqueueRequests(event.requests)
return True
if event.type() == TCPING_CANCEL_EVENT_TYPE:
self.cancelRequests(event.requestIds)
return True
return super().event(event)
def enqueueRequests(self, requests):
"""Queue unique endpoints and start work within the current window."""
if self.stopping:
return
self.pendingRequests.extend(requests)
self.drain()
def drain(self):
"""Construct sockets incrementally inside the networking thread."""
while (
not self.stopping
and self.pendingRequests
and len(self.activeProbes) < self.currentConcurrency
):
request = self.pendingRequests.popleft()
probe = self.probeFactory(request, parent=self)
self.activeProbes[request.requestId] = probe
connectWeakly(
probe.finished,
self,
'handleProbeFinished',
sender=probe,
)
probe.start()
@QtCore.Slot(int, object, bool, bool)
def handleProbeFinished(
self,
requestId: int,
result,
deadlineExpired: bool,
cancelled: bool,
):
"""Release one probe, adapt the window, and post a result to the GUI."""
probe = self.activeProbes.pop(requestId, None)
if probe is None:
return
probe.deleteLater()
if not cancelled:
self.recordOutcome(deadlineExpired)
self.postResult(requestId, result)
self.drain()
def recordOutcome(self, deadlineExpired: bool):
"""Use additive increase and timeout-driven multiplicative decrease."""
if deadlineExpired:
self.currentConcurrency = max(self.currentConcurrency // 2, 1)
self.responsiveCompletionCount = 0
return
self.responsiveCompletionCount += 1
growthThreshold = max(self.currentConcurrency * 2, 4)
if (
self.responsiveCompletionCount >= growthThreshold
and self.currentConcurrency < self.maxConcurrency
):
self.currentConcurrency += 1
self.responsiveCompletionCount = 0
def postResult(self, requestId: int, result):
"""Post one immutable result without retaining the GUI scheduler."""
receiver = self.resultReceiver()
if receiver is None:
return
try:
QtCore.QCoreApplication.postEvent(
receiver, TcpingResultEvent(requestId, result)
)
except RuntimeError:
# The receiver's native QObject was deleted between weak resolution
# and the thread-safe event post.
pass
def cancelRequests(self, requestIds):
"""Remove queued requests and abort matching active sockets."""
requestIds = set(requestIds)
self.pendingRequests = collections.deque(
request
for request in self.pendingRequests
if request.requestId not in requestIds
)
for requestId in requestIds:
probe = self.activeProbes.get(requestId)
if probe is not None:
probe.cancel()
self.drain()
@QtCore.Slot()
def shutdown(self):
"""Abort every owned socket before the dedicated event loop stops."""
if self.stopping:
return
self.stopping = True
self.pendingRequests.clear()
for probe in list(self.activeProbes.values()):
probe.cancel()
class TcpingThread(QtCore.QThread):
"""Run and finally destroy one Tcping engine in its own native thread."""
def __init__(self, engine: TcpingEngine, parent=None):
"""Initialize the TcpingThread."""
super().__init__(parent)
self.engine = engine
def run(self):
"""Run the event loop, then synchronously destroy its stopped engine."""
try:
self.exec()
finally:
self.engine.shutdown()
deleteQObject(self.engine)
self.engine = None
+18
View File
@@ -51,6 +51,16 @@ from .MetricsHistory import (
MetricsHistory,
)
from .PluginUIManager import PluginNavigationManager, isCoreActive
from .ProfileTesting import (
DownloadSpeedTestOptions,
LatencyTestOptions,
LatencyTestType,
ProfileTestField,
ProfileTestJobState,
ProfileTestManager,
ProfileTestResult,
ProfileTestTarget,
)
from .SubscriptionImporter import (
SubscriptionImportResult,
SubscriptionImportService,
@@ -98,6 +108,14 @@ __all__ = [
'MetricSample',
'MetricsHistory',
'PluginNavigationManager',
'DownloadSpeedTestOptions',
'LatencyTestOptions',
'LatencyTestType',
'ProfileTestField',
'ProfileTestJobState',
'ProfileTestManager',
'ProfileTestResult',
'ProfileTestTarget',
'SubscriptionImportResult',
'SubscriptionImportService',
'SubscriptionSource',
+4 -1
View File
@@ -9,10 +9,13 @@
- Qt models must bracket live-collection mutations with matching begin/end notifications and keep stored row/index flags
synchronized. After sort/filter, map an action from the current proxy index to the source object and use stable IDs;
display text and row position are not identity.
- Models, delegates, headers, menus, actions, spinners, animations, timers, test schedulers, workers, network helpers, and
- Models, delegates, headers, menus, actions, spinners, animations, timers, network helpers, and
reusable editors each need an explicit owner. Persistent widgets construct/connect once; refresh/show changes state
instead of accumulating objects.
- Keep expensive parsing, aggregation, mapping, screen/network/core work off blocking GUI paths or split it into bounded
event-loop units. Reject superseded worker/reply results before mutating a live model.
- `ServerTableView` owns profile-test selection and repaint only. Submit live selected profiles to `ProfileTestManager`,
forward repository/subscription mutation boundaries, and repaint the exact committed latency/speed cell; do not move
workers, queues, networking, concurrency, port allocation, or result mutation back into the widget.
- Verify behavior plus repeated refresh/show/open/close lifetime, sorted/filtered actions, model notification ranges,
cancellation/stale results, hidden-page rendering, and exact cleanup of worker/core/native resources.
+76 -685
View File
@@ -33,8 +33,9 @@ from Furious.Qt import *
from Furious.Qt.Signals import connectWeakly, singleShotWeakly
from Furious.Qt import gettext as _
from Furious.Service import (
CORE_LOG_CATEGORY,
ConnectionManager,
ProfileTestField,
ProfileTestManager,
ProfileTestResult,
SubscriptionManager,
SubscriptionUpdateBatch,
)
@@ -43,15 +44,12 @@ from Furious.Widget.WaitingSpinner import WaitingSpinner
from PySide6 import QtCore
from PySide6.QtGui import *
from PySide6.QtWidgets import *
from PySide6.QtNetwork import *
from typing import Callable, Union
import re
import logging
import icmplib
import functools
import collections
__all__ = ['ServerTableView']
@@ -63,21 +61,6 @@ registerAppSettings('ServerWidgetSectionSizeTable')
registerAppSettings('UserServersHeaderViewState')
def appIsExiting() -> bool:
"""Return the app is exiting value used by the application."""
app = APP()
if app is None:
return True
else:
isExiting = getattr(app, 'isExiting', None)
if callable(isExiting):
return isExiting()
else:
return True
class MBoxUpdateSubsInfo(AppQMessageBox):
"""Represent m box update subs info."""
@@ -148,512 +131,6 @@ class MBoxUpdateSubsInfo(AppQMessageBox):
self.moveToCenter()
class TestPingLatencyWorker(QtCore.QObject, QtCore.QRunnable):
"""Run test ping latency work in the background."""
finished = QtCore.Signal()
def __init__(self, factory: ServerProfile):
# Explictly called __init__
"""Initialize the TestPingLatencyWorker."""
QtCore.QObject.__init__(self)
QtCore.QRunnable.__init__(self)
self.factory = factory
def run(self):
"""Run the test ping latency worker task."""
index = self.factory.index
if self.factory.deleted or index < 0 or index >= len(Storage.UserServers()):
# Invalid item. Do nothing
return
assert isinstance(self.factory, ServerProfile)
try:
result = icmplib.ping(
self.factory.itemAddress,
count=1,
timeout=2,
interval=1,
)
except Exception as ex:
# Any non-exit exceptions
self.factory.metadata.latency = classname(ex)
else:
# Result address should not be empty
if result.address and result.is_alive:
self.factory.metadata.latency = f'{round(result.avg_rtt)}ms'
else:
if result.packet_loss == 1:
self.factory.metadata.latency = 'Timeout'
else:
self.factory.metadata.latency = 'Error'
finally:
# Extra guard
if not appIsExiting():
self.finished.emit()
class TestTcpingLatencyWorker(QtCore.QObject, QtCore.QRunnable):
"""Run test tcping latency work in the background."""
finished = QtCore.Signal()
def __init__(self, factory: ServerProfile):
# Explictly called __init__
"""Initialize the TestTcpingLatencyWorker."""
QtCore.QObject.__init__(self)
QtCore.QRunnable.__init__(self)
self.factory = factory
def run(self):
"""Run the test tcping latency worker task."""
index = self.factory.index
if self.factory.deleted or index < 0 or index >= len(Storage.UserServers()):
# Invalid item. Do nothing
return
assert isinstance(self.factory, ServerProfile)
try:
sent, rtts = tcping(
self.factory.itemAddress,
int(self.factory.itemPort.split(',')[0]),
count=1,
timeout=2,
interval=1,
)
except Exception as ex:
# Any non-exit exceptions
self.factory.metadata.latency = classname(ex)
else:
if rtts:
self.factory.metadata.latency = f'{round(rtts[0] * 1000)}ms'
else:
self.factory.metadata.latency = 'Timeout'
finally:
# Extra guard
if not appIsExiting():
self.finished.emit()
class TestDownloadSpeedWorker(HttpGetManager):
"""Run test download speed work in the background."""
progressed = QtCore.Signal()
finished = QtCore.Signal(object)
def __init__(
self,
factory: ServerProfile,
port: int,
timeout: int,
parent=None,
**kwargs,
):
"""Initialize the TestDownloadSpeedWorker."""
actionMessage = kwargs.pop('actionMessage', 'test download speed')
super().__init__(parent, actionMessage=actionMessage)
self.factory = factory
self.port = port
self.timeout = timeout
self.kwargs = kwargs
self.hasSpeedResult = False
self.totalBytesRead = 0
self.hasDataCounter = 0
self.coreManager = ConnectionManager()
self.networkReply = None
self.elapsedTimer = QtCore.QElapsedTimer()
self.timeoutTimer = QtCore.QTimer(self)
self.timeoutTimer.setSingleShot(True)
connectWeakly(
self.timeoutTimer.timeout,
self,
'handleTimeout',
sender=self.timeoutTimer,
)
def completionCallback(self, **kwargs):
"""Perform the required completion hook."""
self.timeoutTimer.stop()
self.finished.emit(self)
def sync(self):
# Extra guard
"""Persist the current test download speed worker data."""
if not appIsExiting():
self.progressed.emit()
def isFinished(self) -> bool:
"""Return whether finished."""
if isinstance(self.networkReply, QNetworkReply):
return self.networkReply.isFinished()
else:
return True
def abort(self):
"""Cancel the active test download speed worker operation."""
if isinstance(self.networkReply, QNetworkReply):
self.networkReply.abort()
def handleTimeout(self):
"""Handle timeout."""
try:
if not self.isFinished():
self.abort()
finally:
self.runCompletionCallback()
def coreExitCallback(self, config: CoreConfiguration, exitcode: int):
"""Handle the core exit callback."""
try:
if exitcode == CoreRuntime.ExitCode.ConfigurationError.value:
self.factory.metadata.speed = 'Invalid'
self.sync()
elif exitcode == CoreRuntime.ExitCode.ServerStartFailure.value:
self.factory.metadata.speed = 'Core start failed'
self.sync()
elif exitcode == CoreRuntime.ExitCode.SystemShuttingDown.value:
pass
else:
self.factory.metadata.speed = f'Core exited {exitcode}'
self.sync()
finally:
self.runCompletionCallback()
def _startCoreRuntime(self, config) -> bool:
"""Prepare and start a download test through its runtime factory."""
configcopy = getPluginRegistry().prepareDownloadTest(config, self.port)
if configcopy is None:
self.factory.metadata.speed = 'Invalid'
self.sync()
return False
self.factory.metadata.speed = 'Starting'
self.sync()
return self.coreManager.start(
configcopy,
AppBuiltinRouting.Global.value,
self.coreExitCallback,
msgCallbackCore=AppLogManager().callback(CORE_LOG_CATEGORY),
deepcopy=False,
proxyModeOnly=True,
log=False,
)
def start(self):
"""Start the test download speed worker."""
try:
if appIsExiting():
raise
index = self.factory.index
if self.factory.deleted or index < 0 or index >= len(Storage.UserServers()):
# Invalid item. Do nothing
return
assert isinstance(self.factory, ServerProfile)
if not self.factory.isValid():
# Configuration is invalid
self.factory.metadata.speed = 'Invalid'
self.sync()
else:
if not self._startCoreRuntime(self.factory) or appIsExiting():
return
self.configureHttpProxy(f'127.0.0.1:{self.port}')
# Use custom network speed test URL if possible
settings = AppSettings.get('CustomNetworkSpeedTestURL')
if isinstance(settings, str):
url = settings
else:
url = NETWORK_SPEED_TEST_URL
self.networkReply = self.webGET(url, **self.kwargs)
self.elapsedTimer.start()
self.timeoutTimer.start(self.timeout)
finally:
if self.networkReply is None:
self.runCompletionCallback()
def successCallback(self, networkReply, **kwargs):
"""Handle a successful network operation."""
if self.coreManager.allRunning():
self.totalBytesRead += networkReply.readAll().length()
# Convert to seconds
elapsedSecond = self.elapsedTimer.elapsed() / 1000
downloadSpeed = self.totalBytesRead / elapsedSecond / 1024 / 1024
self.factory.metadata.speed = f'{downloadSpeed:.2f} MiB/s'
else:
self.factory.metadata.speed = 'Core start failed'
self.coreManager.stopAll()
self.sync()
def hasDataCallback(self, networkReply, **kwargs):
"""Handle newly available network response data."""
self.hasDataCounter += 1
if self.coreManager.allRunning():
self.totalBytesRead += networkReply.readAll().length()
# Convert to seconds
elapsedSecond = self.elapsedTimer.elapsed() / 1000
downloadSpeed = self.totalBytesRead / elapsedSecond / 1024 / 1024
# Has speed test result
self.hasSpeedResult = True
self.factory.metadata.speed = f'{downloadSpeed:.2f} MiB/s'
# Limited to save CPU resources
if self.hasDataCounter % 25 == 0:
self.sync()
def failureCallback(self, networkReply, **kwargs):
"""Handle a failed network operation."""
if not self.hasSpeedResult:
if not self.coreManager.allRunning():
# Core ExitCallback has been called
return
if (
networkReply.error()
== QNetworkReply.NetworkError.OperationCanceledError
):
# Canceled by application
self.factory.metadata.speed = 'Canceled'
else:
try:
error = networkReply.error().name
except Exception:
# Any non-exit exceptions
error = 'UnknownError'
if isinstance(error, bytes):
# Some old version PySide6 returns it as bytes. Protect it.
error = error.decode('utf-8', 'replace')
elif isinstance(error, str):
pass
else:
error = 'UnknownError'
if error != 'UnknownError' and error.endswith('Error'):
self.factory.metadata.speed = error[:-5]
else:
self.factory.metadata.speed = error
self.coreManager.stopAll()
self.sync()
class DownloadSpeedTestJob:
"""Represent download speed test job."""
def __init__(
self,
index: int,
factory: ServerProfile,
timeout: int,
logActionMessage=False,
):
"""Initialize the DownloadSpeedTestJob."""
super().__init__()
self.index = index
self.factory = factory
self.timeout = timeout
self.logActionMessage = logActionMessage
class DownloadSpeedTestScheduler(QtCore.QObject):
"""Schedule and coordinate download speed test jobs."""
SinglePort = 20809
MultiPortStart = 30000
MultiPortStop = 40000
def __init__(self, table, isMulti: bool, parent=None):
"""Initialize the DownloadSpeedTestScheduler."""
super().__init__(parent)
self.table = table
self.isMulti = isMulti
self.maxConcurrency = max(OS_CPU_COUNT // 2, 1) if isMulti else 1
self.queue = collections.deque()
self.activeJobs = {}
self.activePorts = set()
self.nextMultiPort = self.MultiPortStart
self.drainScheduled = False
def enqueue(
self,
index: int,
factory: ServerProfile,
timeout: int,
logActionMessage=False,
):
"""Handle enqueue for the download speed test scheduler."""
self.queue.append(
DownloadSpeedTestJob(index, factory, timeout, logActionMessage)
)
self.scheduleDrain()
def enqueueMany(self, jobs: list[DownloadSpeedTestJob]):
"""Handle enqueue many for the download speed test scheduler."""
self.queue.extend(jobs)
self.scheduleDrain()
def cancelAll(self):
"""Return whether cel all."""
self.queue.clear()
for worker, _, _ in list(self.activeJobs.values()):
assert isinstance(worker, TestDownloadSpeedWorker)
if not worker.isFinished():
worker.abort()
worker.coreManager.stopAll()
worker.runCompletionCallback()
def scheduleDrain(self):
"""Handle schedule drain for the download speed test scheduler."""
if self.drainScheduled:
return
self.drainScheduled = True
singleShotWeakly(0, self, 'drain')
def drain(self):
"""Handle drain for the download speed test scheduler."""
self.drainScheduled = False
if appIsExiting():
self.cancelAll()
return
while self.queue and len(self.activeJobs) < self.maxConcurrency:
job = self.queue.popleft()
assert isinstance(job.factory, ServerProfile)
if job.factory.deleted:
continue
port = self.allocatePort()
if port is None:
self.queue.appendleft(job)
break
self.startJob(job, port)
def allocatePort(self) -> Union[int, None]:
"""Return whether allocate port."""
if not self.isMulti:
if self.activeJobs:
return None
return self.SinglePort
portRange = self.MultiPortStop - self.MultiPortStart
for _ in range(portRange):
port = self.nextMultiPort
self.nextMultiPort += 1
if self.nextMultiPort >= self.MultiPortStop:
self.nextMultiPort = self.MultiPortStart
if port not in self.activePorts:
self.activePorts.add(port)
return port
return None
def releasePort(self, port: int):
"""Handle release port for the download speed test scheduler."""
self.activePorts.discard(port)
def startJob(self, job: DownloadSpeedTestJob, port: int):
"""Start job."""
worker = TestDownloadSpeedWorker(
job.factory,
port,
job.timeout,
parent=self,
logActionMessage=job.logActionMessage,
)
worker.progressed.connect(
functools.partial(
self.table.flushDownloadSpeedItem,
job.index,
job.factory,
)
)
self.activeJobs[id(worker)] = (worker, job, port)
connectWeakly(
worker.finished,
self,
'handleWorkerFinished',
sender=worker,
)
worker.start()
@QtCore.Slot(object)
def handleWorkerFinished(self, worker):
"""Handle worker finished."""
workerId = id(worker)
try:
_, _, port = self.activeJobs.pop(workerId)
except KeyError:
return
self.releasePort(port)
# Completed workers are children of the long-lived scheduler. Merely
# removing the Python dictionary entry would leave every worker (and
# its network/timer children) in the scheduler's QObject tree.
worker.deleteLater()
self.scheduleDrain()
class DeleteServersProgressDialog(AppQTransientDialog):
"""Present progress and cancellation controls for delete servers."""
@@ -791,6 +268,7 @@ class DeleteServersProgressDialog(AppQTransientDialog):
Storage.UserServers().pop(deleteIndex)
self.table.sourceModel.endRemoveRows()
self.table.reconcileProfileTestJobs()
if not self.deletedActivated and deleteIndex < Storage.UserActivatedItemIndex():
AppSettings.set(
@@ -1306,19 +784,20 @@ class ServerTableView(
self.subsManager = SubscriptionManager(parent=self)
self.subsManager.subscriptionsChanged.connect(self._handleSubscriptionsChanged)
self.subsManager.subscriptionCommitted.connect(
self._handleSubscriptionCommitted
)
self.subsManager.updateCompleted.connect(
self._handleSubscriptionUpdateCompleted
)
self.downloadSpeedScheduler = DownloadSpeedTestScheduler(
self.profileTestManager = ProfileTestManager(parent=self)
connectWeakly(
self.profileTestManager.resultApplied,
self,
isMulti=False,
parent=self,
)
self.downloadSpeedMultiScheduler = DownloadSpeedTestScheduler(
self,
isMulti=True,
parent=self,
'_handleProfileTestResultApplied',
sender=self.profileTestManager,
)
self.configurationEditor = configurationEditorFactory()
@@ -1933,6 +1412,11 @@ class ServerTableView(
def flushRow(self, row: int, item: ServerProfile):
"""Refresh row."""
# flushRow is the established commit notification for both structured
# and JSON profile editors. Reconcile before presenting a replacement
# or in-place connection mutation.
self.reconcileProfileTestJobs()
itemIndex = item.index
if item.deleted or itemIndex < 0 or itemIndex >= len(Storage.UserServers()):
@@ -2254,6 +1738,8 @@ class ServerTableView(
'ActivatedItemIndex', str(Storage.UserActivatedItemIndex() - 1)
)
self.reconcileProfileTestJobs()
# Refresh index
self.sourceModel.refreshIndexes()
self.sourceModel.emitAllChanged()
@@ -2358,185 +1844,80 @@ class ServerTableView(
self.setCurrentIndex(activatedItem)
self.scrollTo(activatedItem)
def rowFromFactory(self, fallbackIndex: int, factory: ServerProfile) -> int:
"""Return the row from factory value."""
if (
0 <= factory.index < len(Storage.UserServers())
and Storage.UserServers()[factory.index] is factory
):
return factory.index
@QtCore.Slot(object, object)
def _handleProfileTestResultApplied(
self,
profile: ServerProfile,
result: ProfileTestResult,
):
"""Repaint the one repository cell committed by the test service."""
profiles = Storage.UserServers()
row = profile.index
if (
0 <= fallbackIndex < len(Storage.UserServers())
and Storage.UserServers()[fallbackIndex] is factory
):
return fallbackIndex
for index, item in enumerate(Storage.UserServers()):
if item is factory:
return index
return -1
def flushDownloadSpeedItem(self, fallbackIndex: int, factory: ServerProfile):
"""Refresh download speed item."""
index = self.rowFromFactory(fallbackIndex, factory)
if index < 0:
if row < 0 or row >= len(profiles) or profiles[row] is not profile:
return
self.flushItem(index, self.Headers.index('Speed'), factory)
column = 'Latency' if result.field is ProfileTestField.Latency else 'Speed'
self.flushItem(row, self.Headers.index(column), profile)
def _selectedProfilesForTesting(self):
"""Resolve the current row selection into live repository profiles."""
profiles = Storage.UserServers()
return [
profiles[index]
for index in self.selectedIndex
if 0 <= index < len(profiles)
]
def reconcileProfileTestJobs(self):
"""Notify the test service after a live profile collection mutation."""
self.profileTestManager.reconcileProfiles()
def testSelectedItemPingLatency(self):
"""Handle test selected item ping latency for the user servers Qt table view."""
indexes = self.selectedIndex
if len(indexes) == 0:
# Nothing selected. Do nothing
return
# Real selected factory
references = list(Storage.UserServers()[index] for index in indexes)
for index, reference in zip(indexes, references):
if appIsExiting():
break
assert isinstance(reference, ServerProfile)
if reference.deleted:
continue
worker = TestPingLatencyWorker(reference)
worker.setAutoDelete(True)
worker.finished.connect(
functools.partial(
self.flushItem,
index,
self.Headers.index('Latency'),
reference,
)
)
AppThreadPool().start(worker)
"""Request ICMP latency tests for selected profiles."""
self.profileTestManager.testPing(self._selectedProfilesForTesting())
def testSelectedItemTcpingLatency(self):
"""Handle test selected item tcping latency for the user servers Qt table view."""
indexes = self.selectedIndex
if len(indexes) == 0:
# Nothing selected. Do nothing
return
# Real selected factory
references = list(Storage.UserServers()[index] for index in indexes)
for index, reference in zip(indexes, references):
if appIsExiting():
break
assert isinstance(reference, ServerProfile)
if reference.deleted:
continue
worker = TestTcpingLatencyWorker(reference)
worker.setAutoDelete(True)
worker.finished.connect(
functools.partial(
self.flushItem,
index,
self.Headers.index('Latency'),
reference,
)
)
AppThreadPool().start(worker)
def testDownloadSpeedByFactory(
self,
index: int,
factory: ServerProfile,
port: int,
timeout: int,
isMulti: bool,
counter=0,
step=100,
logActionMessage=False,
):
"""Handle test download speed by factory for the user servers Qt table view."""
scheduler = (
self.downloadSpeedMultiScheduler if isMulti else self.downloadSpeedScheduler
)
scheduler.enqueue(index, factory, timeout, logActionMessage)
def testSelectedItemDownloadSpeedWithTimeoutXXX(
self,
scheduler: DownloadSpeedTestScheduler,
timeout: int,
):
"""Handle test selected item download speed with timeout xxx for the user servers Qt table view."""
indexes = self.selectedIndex
if len(indexes) == 0:
# Nothing selected. Do nothing
return
# Real selected factory
references = list(Storage.UserServers()[index] for index in indexes)
jobs = list()
for index, reference in zip(indexes, references):
jobs.append(DownloadSpeedTestJob(index, reference, timeout))
scheduler.enqueueMany(jobs)
"""Request asynchronous TCP latency tests for selected profiles."""
self.profileTestManager.testTcping(self._selectedProfilesForTesting())
def testSelectedItemDownloadSpeedWithTimeout(self, timeout: int):
"""Handle test selected item download speed with timeout for the user servers Qt table view."""
self.testSelectedItemDownloadSpeedWithTimeoutXXX(
self.downloadSpeedScheduler,
timeout,
"""Request serial download tests for selected profiles."""
self.profileTestManager.testDownloadSpeed(
self._selectedProfilesForTesting(),
timeoutMilliseconds=timeout,
concurrent=False,
)
def testSelectedItemDownloadSpeedWithTimeoutMulti(self, timeout: int):
"""Handle test selected item download speed with timeout multi for the user servers Qt table view."""
self.testSelectedItemDownloadSpeedWithTimeoutXXX(
self.downloadSpeedMultiScheduler,
timeout,
"""Request concurrent download tests for selected profiles."""
self.profileTestManager.testDownloadSpeed(
self._selectedProfilesForTesting(),
timeoutMilliseconds=timeout,
concurrent=True,
)
def testSelectedItemDownloadSpeed(self):
"""Run the retained single-threaded download-test API for selected rows.
"""Run the retained serial download-test API for selected rows.
The Home context menu uses the multithreaded variant; this method remains
The Home context menu uses the concurrent variant; this method remains
available for programmatic callers that explicitly need serial scheduling.
"""
self.testSelectedItemDownloadSpeedWithTimeout(5000)
def testSelectedItemDownloadSpeedMulti(self):
"""Handle test selected item download speed multi for the user servers Qt table view."""
"""Request concurrent download tests for selected profiles."""
self.testSelectedItemDownloadSpeedWithTimeoutMulti(5000)
def clearSelectedItemTestResult(self):
"""Clear selected item test result."""
indexes = self.selectedIndex
if len(indexes) == 0:
# Nothing selected. Do nothing
return
for index in indexes:
factory = Storage.UserServers()[index]
factory.metadata.latency = ''
factory.metadata.speed = ''
self.flushItem(index, self.Headers.index('Latency'), factory)
self.flushItem(index, self.Headers.index('Speed'), factory)
self.profileTestManager.clearResults(self._selectedProfilesForTesting())
def cleanup(self):
"""Release resources owned by the user servers Qt table view."""
self.downloadSpeedScheduler.cancelAll()
self.downloadSpeedMultiScheduler.cancelAll()
self.profileTestManager.shutdown()
self._clearSubscriptionActions()
def updateSubsByUnique(self, unique: str, httpProxy: Union[str, None], **kwargs):
@@ -2558,6 +1939,8 @@ class ServerTableView(
@QtCore.Slot()
def _handleSubscriptionsChanged(self):
"""Refresh the table after the service commits repository changes."""
self.reconcileProfileTestJobs()
self.sourceModel.beginResetModel()
self.sourceModel.endResetModel()
self.sourceModel.refreshIndexes()
@@ -2568,6 +1951,14 @@ class ServerTableView(
self.activeServerChanged.emit()
@QtCore.Slot(str)
def _handleSubscriptionCommitted(self, unique: str):
"""Clear results and invalidate only work owned by one committed group."""
self.profileTestManager.invalidateSubscriptions(
{unique},
clearResults=True,
)
@QtCore.Slot(object)
def _handleSubscriptionUpdateCompleted(self, batch: SubscriptionUpdateBatch):
"""Present one semantic subscription update result batch."""
+3 -2
View File
@@ -45,6 +45,7 @@ strategy in an individual test.
| Controller state and error transitions with injected runtimes | `test_controllers.py` |
| SOCKS and SIP002 Shadowsocks codecs and generated round trips | `test_socks_uri.py`, `test_shadowsocks_uri.py` |
| Subscription workflow, timers, stale requests, and reconciliation | `test_subscription_manager.py`, `test_subscription_sync.py` |
| Service-first profile-test identity, explicit results, endpoint deduplication, adaptive Tcping, cancellation, late callbacks, and worker/thread lifetime | `test_profile_test_jobs.py` |
| External process launch, output, shutdown, threads, TUN metadata | `test_external_core.py` |
| Backend structured-editor observational load and unknown-value preservation | `test_backend_editor_contract.py` |
| Xray asset checksum validation and atomic replacement | `test_xray_asset_download.py` |
@@ -108,7 +109,7 @@ Then run the desired test tier.
python -m unittest discover -s tests -v
# Regular logic, persistence, plugin, controller, codec, and UI regressions
python -m unittest tests.test_interface tests.test_models_and_services tests.test_repository_contracts tests.test_architecture_refactors tests.test_plugin_architecture tests.test_hysteria1_protocol tests.test_hysteria2_compatibility tests.test_controllers tests.test_subscription_manager tests.test_subscription_sync tests.test_socks_uri tests.test_shadowsocks_uri tests.test_backend_editor_contract tests.test_xray_asset_download tests.test_native_tun_semantics tests.test_metrics_behavior tests.test_endpoint_info tests.test_service_runtime tests.test_frozenlib tests.test_isolation_and_navigation tests.test_main_window_geometry tests.test_dialog_geometry tests.test_ui_behavior tests.test_qt_interactions tests.test_stylesheet_states tests.test_theme_transition tests.test_public_api -v
python -m unittest tests.test_interface tests.test_models_and_services tests.test_repository_contracts tests.test_architecture_refactors tests.test_plugin_architecture tests.test_hysteria1_protocol tests.test_hysteria2_compatibility tests.test_controllers tests.test_subscription_manager tests.test_subscription_sync tests.test_profile_test_jobs tests.test_socks_uri tests.test_shadowsocks_uri tests.test_backend_editor_contract tests.test_xray_asset_download tests.test_native_tun_semantics tests.test_metrics_behavior tests.test_endpoint_info tests.test_service_runtime tests.test_frozenlib tests.test_isolation_and_navigation tests.test_main_window_geometry tests.test_dialog_geometry tests.test_ui_behavior tests.test_qt_interactions tests.test_stylesheet_states tests.test_theme_transition tests.test_public_api -v
# Direct Qt/process integration and destruction/lifetime checks
python -m unittest tests.test_application_process tests.test_external_core tests.test_layout_matrix tests.test_qt_lifetime -v
@@ -125,7 +126,7 @@ python -m unittest tests.test_very_heavy -v
python -m unittest tests.test_log_manager_generation.VeryHeavyGenerationLogManagerTest -v
# Shared-state order-independence spot check
python -m unittest tests.test_public_api tests.test_theme_transition tests.test_stylesheet_states tests.test_qt_interactions tests.test_ui_behavior tests.test_dialog_geometry tests.test_main_window_geometry tests.test_isolation_and_navigation tests.test_frozenlib tests.test_service_runtime tests.test_endpoint_info tests.test_metrics_behavior tests.test_native_tun_semantics tests.test_xray_asset_download tests.test_backend_editor_contract tests.test_shadowsocks_uri tests.test_socks_uri tests.test_subscription_sync tests.test_subscription_manager tests.test_controllers tests.test_hysteria2_compatibility tests.test_hysteria1_protocol tests.test_plugin_architecture tests.test_architecture_refactors tests.test_repository_contracts tests.test_models_and_services tests.test_interface -v
python -m unittest tests.test_public_api tests.test_theme_transition tests.test_stylesheet_states tests.test_qt_interactions tests.test_ui_behavior tests.test_dialog_geometry tests.test_main_window_geometry tests.test_isolation_and_navigation tests.test_frozenlib tests.test_service_runtime tests.test_endpoint_info tests.test_metrics_behavior tests.test_native_tun_semantics tests.test_xray_asset_download tests.test_backend_editor_contract tests.test_shadowsocks_uri tests.test_socks_uri tests.test_profile_test_jobs tests.test_subscription_sync tests.test_subscription_manager tests.test_controllers tests.test_hysteria2_compatibility tests.test_hysteria1_protocol tests.test_plugin_architecture tests.test_architecture_refactors tests.test_repository_contracts tests.test_models_and_services tests.test_interface -v
python -m unittest discover -s tests -v
```
+926
View File
@@ -0,0 +1,926 @@
# 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/>.
"""Protect profile-test identity, scheduling, cancellation, and Qt lifetimes."""
from __future__ import annotations
from Furious.Backends.Configuration import ConfigXray
from Furious.Frozenlib import AppSettings, OS_CPU_COUNT
from Furious.Models import ServerProfile
from Furious.Repository import Storage
from Furious.Service.ProfileTesting import (
DownloadSpeedTestOptions,
LatencyTestOptions,
LatencyTestType,
ProfileTestField,
ProfileTestJobState,
ProfileTestManager,
ProfileTestResult,
ProfileTestTarget,
_DownloadSpeedWorker,
_PingResultEvent,
)
from Furious.Service.TcpingService import TcpingEngine
from Furious.Widget.ServerTableView import ServerTableView
from PySide6 import QtCore, QtNetwork
from PySide6.QtWidgets import QWidget
from shiboken6 import isValid
from tests.support import (
application,
collectAtBoundary,
isolatedSettings,
processQtEvents,
waitFor,
)
import threading
import unittest
from types import SimpleNamespace
from unittest import mock
class _ControlledDownloadWorker(QtCore.QObject):
"""Expose deterministic progress and completion through real Qt signals."""
progressed = QtCore.Signal(object, object)
finished = QtCore.Signal(object, object)
instances = []
def __init__(self, profile, port, options, parent=None):
"""Retain explicit inputs without changing the profile snapshot."""
super().__init__(parent)
self.profile = profile
self.port = port
self.options = options
self.cancelCount = 0
self.terminal = False
self.result = ProfileTestResult(
ProfileTestField.DownloadSpeed,
'',
terminal=False,
)
self.instances.append(self)
def start(self):
"""Leave completion under explicit test control."""
def publish(self, speed):
"""Publish one non-terminal result without mutating the snapshot."""
if self.terminal:
return
self.result = ProfileTestResult(
ProfileTestField.DownloadSpeed,
speed,
terminal=False,
)
self.progressed.emit(self, self.result)
def finish(self, speed):
"""Complete this worker exactly once with an explicit result."""
if self.terminal:
return
self.result = ProfileTestResult(ProfileTestField.DownloadSpeed, speed)
self.terminal = True
self.finished.emit(self, self.result)
def cancel(self):
"""Record cancellation and publish one terminal cancellation signal."""
if self.terminal:
return
self.cancelCount += 1
self.terminal = True
self.finished.emit(self, ProfileTestResult(ProfileTestField.DownloadSpeed, ''))
class _ControlledLatencyWorker(QtCore.QRunnable):
"""Leave a Ping QRunnable's result ordering under test control."""
def __init__(self, job, receiver):
"""Retain the explicit job and scheduler receiver."""
super().__init__()
self.job = job
self.receiver = receiver
def run(self):
"""Do nothing because the controlled pool does not execute workers."""
def finish(self, result):
"""Post one result through the production scheduler event boundary."""
QtCore.QCoreApplication.postEvent(
self.receiver,
_PingResultEvent(self.job, result),
)
class _ControlledThreadPool:
"""Collect scheduler-started QRunnables without executing them."""
def __init__(self):
"""Initialize an empty worker list."""
self.started = []
def start(self, worker):
"""Record one worker for explicit signal delivery."""
self.started.append(worker)
@staticmethod
def waitForDone(_timeout):
"""Represent a deterministic pool with no background execution."""
return True
class _CoreManagerProbe:
"""Record exact runtime disposal before worker QObject deletion."""
def __init__(self):
"""Initialize a stopped runtime owner."""
self.stopCount = 0
@staticmethod
def allRunning():
"""Represent the core-death side of the reported race."""
return False
def stopAll(self):
"""Record callback and runtime disposal."""
self.stopCount += 1
class _ImmediateCoreManager(_CoreManagerProbe):
"""Launch immediately while recording the embedded-core wait policy."""
def __init__(self):
"""Initialize one ready runtime and its captured start calls."""
super().__init__()
self.startCalls = []
def start(self, *args, **kwargs):
"""Record a non-blocking launch and report immediate process creation."""
self.startCalls.append((args, kwargs))
return True
@staticmethod
def allRunning():
"""Keep the fake runtime alive until the worker is cancelled."""
return True
class _CancelDuringStartDownloadWorker(_DownloadSpeedWorker):
"""Re-enter subscription invalidation while a worker is starting."""
instances = []
def __init__(self, *args, **kwargs):
"""Retain wrappers so native deletion remains observable."""
super().__init__(*args, **kwargs)
self.instances.append(self)
def _startCoreRuntime(self):
"""Cancel this worker at the reentrant runtime-start boundary."""
self.parent().invalidateSubscriptions(
{self.profile.metadata.subscriptionSource}
)
processQtEvents()
return True
class ProfileTestServiceTest(unittest.TestCase):
"""Exercise the self-contained profile-test subsystem."""
@classmethod
def setUpClass(cls):
"""Create the shared isolated QApplication."""
application()
def setUp(self):
"""Reset controlled owners for each case."""
self.profiles = []
self.managers = []
_ControlledDownloadWorker.instances = []
_CancelDuringStartDownloadWorker.instances = []
def tearDown(self):
"""Stop service-owned threads and collect deferred QObjects."""
for manager in reversed(self.managers):
manager.shutdown()
manager.deleteLater()
application().threadPool.waitForDone(5000)
collectAtBoundary()
@staticmethod
def _profile(name: str, address: str, port: int = 443):
"""Build one valid-enough profile with stable identity."""
return ServerProfile.fromConfiguration(
ConfigXray(
{
'outbounds': [
{
'tag': 'proxy',
'protocol': 'vless',
'settings': {
'vnext': [
{
'address': address,
'port': port,
'users': [{}],
}
]
},
}
]
}
),
{'displayName': name},
)
def _setProfiles(self, profiles):
"""Commit one current repository view for the injected provider."""
self.profiles[:] = profiles
for index, profile in enumerate(self.profiles):
profile.index = index
def _manager(self, profiles, *, controlledDownloads=True):
"""Construct one manager around an injected current-profile provider."""
self._setProfiles(profiles)
manager = ProfileTestManager(
profilesProvider=lambda: self.profiles,
pingConcurrency=1,
tcpingConcurrency=2,
downloadConcurrency=1,
)
if controlledDownloads:
manager._serialDownloadScheduler.workerFactory = _ControlledDownloadWorker
manager._concurrentDownloadScheduler.workerFactory = (
_ControlledDownloadWorker
)
self.managers.append(manager)
return manager
def testDefaultConcurrencyKeepsBlockingPingInPrivateHalfCpuPool(self):
"""Keep blocking Ping off shared workers with the requested default limit."""
manager = ProfileTestManager(profilesProvider=lambda: self.profiles)
self.managers.append(manager)
scheduler = manager._latencyScheduler
self.assertEqual(scheduler.maxConcurrency, max(OS_CPU_COUNT // 2, 1))
self.assertIsInstance(scheduler.threadPool, QtCore.QThreadPool)
self.assertIsNot(scheduler.threadPool, application().threadPool)
self.assertEqual(
scheduler.threadPool.maxThreadCount(),
scheduler.maxConcurrency,
)
def testConcurrentDownloadLaunchesEveryAdmittedCoreBeforeReadinessWait(self):
"""Do not serialize admitted jobs behind each core's readiness grace."""
profiles = [
self._profile(f'profile-{index}', f'{index}.example') for index in range(4)
]
manager = self._manager(profiles, controlledDownloads=False)
scheduler = manager._concurrentDownloadScheduler
scheduler.maxConcurrency = len(profiles)
workers = []
def workerFactory(*args, **kwargs):
"""Create a real worker around an immediate fake core runtime."""
worker = _DownloadSpeedWorker(*args, **kwargs)
worker.CoreStartupGraceMilliseconds = 60_000
worker.coreManager = _ImmediateCoreManager()
workers.append(worker)
return worker
scheduler.workerFactory = workerFactory
registry = mock.Mock()
registry.prepareDownloadTest.side_effect = lambda profile, _port: profile
with (
mock.patch(
'Furious.Service.ProfileTesting.getPluginRegistry',
return_value=registry,
),
mock.patch('Furious.Service.ProfileTesting.AppLogManager'),
):
manager.testDownloadSpeed(profiles, concurrent=True)
processQtEvents()
self.assertEqual(len(workers), len(profiles))
self.assertFalse(scheduler.queue)
self.assertEqual(len(scheduler.activeJobs), len(profiles))
for profile, worker in zip(profiles, workers):
self.assertEqual(profile.metadata.speed, 'Starting')
self.assertIsNone(worker.networkReply)
self.assertTrue(worker.coreStartupTimer.isActive())
self.assertEqual(len(worker.coreManager.startCalls), 1)
self.assertIs(
worker.coreManager.startCalls[0][1].get('waitCore'),
False,
)
manager.shutdown()
self.assertTrue(
all(not worker.coreStartupTimer.isActive() for worker in workers)
)
self.assertTrue(waitFor(lambda: all(not isValid(worker) for worker in workers)))
def testDownloadReconciliationCancelsStaleAndStartsNextValidTarget(self):
"""Compact stale work while preserving a current reordered target."""
active = self._profile('active', 'active.example')
queued = self._profile('queued', 'queued.example')
valid = self._profile('valid', 'valid.example')
manager = self._manager((active, queued, valid))
scheduler = manager._serialDownloadScheduler
manager.testDownloadSpeed(
(active, queued, valid),
concurrent=False,
)
processQtEvents()
first = _ControlledDownloadWorker.instances[0]
active.deleted = True
queued.deleted = True
self._setProfiles((valid,))
manager.reconcileProfiles()
self.assertEqual(first.cancelCount, 1)
self.assertEqual(len(scheduler.queue), 1)
first.publish('stale result')
self.assertEqual(valid.metadata.speed, '')
processQtEvents()
self.assertFalse(isValid(first))
second = _ControlledDownloadWorker.instances[1]
self.assertEqual(
second.profile.metadata.profileId,
valid.metadata.profileId,
)
second.finish('8.50 MiB/s')
processQtEvents()
self.assertEqual(valid.metadata.speed, '8.50 MiB/s')
self.assertFalse(scheduler.queue)
self.assertFalse(scheduler.activeJobs)
def testConnectionEditCancelsRunningAndQueuedSnapshots(self):
"""Invalidate a stable ID when its connection fingerprint changes."""
running = self._profile('running', 'old-running.example')
queued = self._profile('queued', 'old-queued.example')
manager = self._manager((running, queued))
scheduler = manager._serialDownloadScheduler
manager.testDownloadSpeed((running, queued), concurrent=False)
processQtEvents()
worker = _ControlledDownloadWorker.instances[0]
running.connection['address'] = 'new-running.example'
queued.connection['address'] = 'new-queued.example'
manager.reconcileProfiles()
self.assertEqual(worker.cancelCount, 1)
self.assertFalse(scheduler.queue)
processQtEvents()
self.assertFalse(scheduler.activeJobs)
worker.publish('stale result')
self.assertEqual(running.metadata.speed, '')
self.assertEqual(queued.metadata.speed, '')
def testReorderAndSubscriptionMovePreserveLogicalJobs(self):
"""Ignore row and ownership metadata changes for unchanged connections."""
firstProfile = self._profile('first', 'first.example')
secondProfile = self._profile('second', 'second.example')
manager = self._manager((firstProfile, secondProfile))
scheduler = manager._serialDownloadScheduler
manager.testDownloadSpeed(
(firstProfile, secondProfile),
concurrent=False,
)
processQtEvents()
firstWorker = _ControlledDownloadWorker.instances[0]
self._setProfiles((secondProfile, firstProfile))
firstProfile.metadata.subscriptionSource = 'new-subscription'
manager.reconcileProfiles()
self.assertEqual(firstWorker.cancelCount, 0)
self.assertEqual(len(scheduler.queue), 1)
firstWorker.finish('1.00 MiB/s')
processQtEvents()
secondWorker = _ControlledDownloadWorker.instances[1]
secondWorker.finish('2.00 MiB/s')
processQtEvents()
self.assertEqual(firstProfile.metadata.speed, '1.00 MiB/s')
self.assertEqual(secondProfile.metadata.speed, '2.00 MiB/s')
self.assertFalse(scheduler.queue)
self.assertFalse(scheduler.activeJobs)
def testRemovingEveryTargetLeavesDownloadSchedulerEmpty(self):
"""Cancel active work and compact the entire pending collection."""
profiles = [
self._profile('one', 'one.example'),
self._profile('two', 'two.example'),
self._profile('three', 'three.example'),
]
manager = self._manager(profiles)
scheduler = manager._serialDownloadScheduler
manager.testDownloadSpeed(profiles, concurrent=False)
processQtEvents()
jobs = [entry[1] for entry in scheduler.activeJobs.values()]
jobs.extend(scheduler.queue)
worker = _ControlledDownloadWorker.instances[0]
for profile in profiles:
profile.deleted = True
self._setProfiles(())
manager.reconcileProfiles()
processQtEvents()
self.assertEqual(worker.cancelCount, 1)
self.assertTrue(all(job.state is ProfileTestJobState.Cancelled for job in jobs))
self.assertFalse(scheduler.queue)
self.assertFalse(scheduler.activeJobs)
def testPingDropsStaleRunningResultAndQueuedTarget(self):
"""Stale-mark active Ping and remove invalid pending work."""
running = self._profile('running', 'running.example')
queued = self._profile('queued', 'queued.example')
valid = self._profile('valid', 'valid.example')
manager = self._manager((running, queued, valid))
scheduler = manager._latencyScheduler
pool = _ControlledThreadPool()
scheduler.threadPool = pool
scheduler.pingWorkerFactory = _ControlledLatencyWorker
manager.testPing((running, queued, valid))
processQtEvents()
first = pool.started[0]
running.deleted = True
queued.deleted = True
self._setProfiles((valid,))
manager.reconcileProfiles()
self.assertIs(first.job.state, ProfileTestJobState.Cancelled)
self.assertEqual(len(scheduler.queue), 1)
first.finish('1ms')
processQtEvents()
self.assertEqual(valid.metadata.latency, '')
self.assertEqual(len(pool.started), 2)
second = pool.started[1]
second.finish('7ms')
processQtEvents()
self.assertEqual(valid.metadata.latency, '7ms')
self.assertFalse(scheduler.queue)
self.assertFalse(scheduler.activeJobs)
def testRealThreadPoolPingDiscardsResultAfterInPlaceEdit(self):
"""Exercise blocking Ping delivery while the target changes mid-call."""
profile = self._profile('profile', 'old.example')
manager = self._manager((profile,))
scheduler = manager._latencyScheduler
started = threading.Event()
release = threading.Event()
def ping(*_args, **_kwargs):
"""Hold the real worker until its target becomes stale."""
started.set()
release.wait(5)
return SimpleNamespace(
address='old.example',
is_alive=True,
avg_rtt=3.2,
packet_loss=0,
)
try:
with mock.patch(
'Furious.Service.ProfileTesting.icmplib.ping',
side_effect=ping,
):
manager.testPing((profile,))
processQtEvents()
self.assertTrue(started.wait(2))
profile.connection['address'] = 'new.example'
manager.reconcileProfiles()
release.set()
self.assertTrue(waitFor(lambda: not scheduler.activeJobs, timeout=5000))
self.assertEqual(profile.metadata.latency, '')
finally:
release.set()
def testTcpingUsesOneNetworkThreadAndDeduplicatesEndpoints(self):
"""Probe a shared endpoint once in the dedicated Qt network thread."""
server = QtNetwork.QTcpServer()
self.assertTrue(
server.listen(QtNetwork.QHostAddress.SpecialAddress.LocalHost, 0)
)
first = self._profile('first', '127.0.0.1', server.serverPort())
second = self._profile('second', '127.0.0.1', server.serverPort())
manager = self._manager((first, second))
scheduler = manager._latencyScheduler
guiEventDelivered = threading.Event()
try:
manager.testTcping((first, second))
QtCore.QTimer.singleShot(0, guiEventDelivered.set)
self.assertEqual(len(scheduler.tcpingRequests), 1)
self.assertEqual(
len(next(iter(scheduler.tcpingRequests.values())).jobs),
2,
)
self.assertIs(scheduler.tcpingEngine.thread(), scheduler.tcpingThread)
self.assertIsNot(scheduler.tcpingThread, application().thread())
self.assertTrue(waitFor(lambda: not scheduler.tcpingRequests, timeout=2000))
self.assertTrue(guiEventDelivered.is_set())
self.assertRegex(first.metadata.latency, r'^\d+ms$')
self.assertEqual(second.metadata.latency, first.metadata.latency)
connections = 0
while server.hasPendingConnections():
socket = server.nextPendingConnection()
connections += 1
socket.deleteLater()
self.assertEqual(connections, 1)
finally:
server.close()
def testTcpingDoesNotCoalesceDifferentTimeoutPolicies(self):
"""Keep endpoint identity separate from explicit execution options."""
first = self._profile('first', '192.0.2.1', 9)
second = self._profile('second', '192.0.2.1', 9)
manager = self._manager((first, second))
scheduler = manager._latencyScheduler
eventSink = QtCore.QObject()
scheduler.tcpingEngine = eventSink
try:
manager.testTcping((first,), timeoutMilliseconds=1000)
manager.testTcping((second,), timeoutMilliseconds=2000)
self.assertEqual(len(scheduler.tcpingRequests), 2)
self.assertEqual(
{
group.request.timeoutMilliseconds
for group in scheduler.tcpingRequests.values()
},
{1000, 2000},
)
finally:
scheduler.tcpingEngine = None
eventSink.deleteLater()
def testTcpingCancellationRemovesEndpointRequestImmediately(self):
"""Invalidate a TCPing group before a network result reaches profiles."""
profile = self._profile('profile', '192.0.2.1', 9)
profile.metadata.subscriptionSource = 'group-a'
manager = self._manager((profile,))
scheduler = manager._latencyScheduler
manager.testTcping((profile,))
group = next(iter(scheduler.tcpingRequests.values()))
job = group.jobs[0]
scheduler.invalidateSubscriptions({'group-a'})
self.assertIs(job.state, ProfileTestJobState.Cancelled)
self.assertFalse(scheduler.tcpingRequests)
self.assertFalse(scheduler.tcpingEndpointRequests)
processQtEvents()
self.assertEqual(profile.metadata.latency, '')
def testTcpingSharedResultFanOutIsBoundedPerGuiBatch(self):
"""Yield between fixed-size write-back batches for a shared endpoint."""
profiles = tuple(
self._profile(f'profile {index}', 'shared.example') for index in range(130)
)
manager = self._manager(profiles)
scheduler = manager._latencyScheduler
eventSink = QtCore.QObject()
scheduler.tcpingEngine = eventSink
try:
manager.testTcping(profiles)
requestId = next(iter(scheduler.tcpingRequests))
scheduler.handleTcpingResult(requestId, '5ms')
scheduler.drainTcpingResults()
self.assertEqual(
sum(profile.metadata.latency == '5ms' for profile in profiles),
scheduler.TcpingResultBatchSize,
)
self.assertTrue(scheduler.tcpingCompletionQueue)
while scheduler.tcpingCompletionQueue:
scheduler.drainTcpingResults()
self.assertTrue(
all(profile.metadata.latency == '5ms' for profile in profiles)
)
finally:
scheduler.tcpingEngine = None
eventSink.deleteLater()
def testTcpingAdaptiveWindowStaysBoundedAndBacksOffOnDeadline(self):
"""Grow conservatively on success and back off after a deadline."""
receiver = QtCore.QObject()
engine = TcpingEngine(receiver, maxConcurrency=8)
self.assertEqual(engine.currentConcurrency, 2)
for _index in range(4):
engine.recordOutcome(False)
self.assertEqual(engine.currentConcurrency, 3)
engine.recordOutcome(True)
self.assertEqual(engine.currentConcurrency, 1)
for _index in range(200):
engine.recordOutcome(False)
self.assertEqual(engine.currentConcurrency, 8)
engine.deleteLater()
receiver.deleteLater()
processQtEvents()
def testTcpingNetworkingThreadHasRepeatableTerminalCleanup(self):
"""Stop and destroy the reusable engine and thread repeatedly."""
for iteration in range(10):
profile = self._profile(
f'profile {iteration}',
'192.0.2.1',
9,
)
manager = self._manager((profile,))
scheduler = manager._latencyScheduler
engine = scheduler.ensureTcpingEngine()
thread = scheduler.tcpingThread
manager.shutdown()
self.assertFalse(thread.isRunning())
manager.deleteLater()
self.managers.remove(manager)
processQtEvents()
self.assertFalse(isValid(engine))
self.assertFalse(isValid(thread))
def testSubscriptionInvalidationCannotDeleteWorkerDuringStart(self):
"""Defer terminal deletion until a reentrant start call unwinds."""
profile = self._profile('profile', 'profile.example')
profile.metadata.subscriptionSource = 'subscription-a'
manager = self._manager((profile,), controlledDownloads=False)
scheduler = manager._serialDownloadScheduler
scheduler.workerFactory = _CancelDuringStartDownloadWorker
for _index in range(25):
manager.testDownloadSpeed((profile,), concurrent=False)
processQtEvents()
worker = _CancelDuringStartDownloadWorker.instances[-1]
self.assertFalse(scheduler.activeJobs)
self.assertFalse(scheduler.activePorts)
self.assertFalse(isValid(worker))
def testSubscriptionCommitInvalidatesOnlyItsOldProfileJobs(self):
"""Clear one committed group while preserving unrelated test work."""
retained = self._profile('retained', 'retained.example')
removed = self._profile('removed', 'removed.example')
replacement = self._profile('replacement', 'replacement.example')
other = self._profile('other', 'other.example')
manual = self._profile('manual', 'manual.example')
for profile in (retained, removed, replacement):
profile.metadata.subscriptionSource = 'group-a'
other.metadata.subscriptionSource = 'group-b'
for profile in (retained, replacement, other, manual):
profile.metadata.latency = f'{profile.itemRemark} latency'
profile.metadata.speed = f'{profile.itemRemark} speed'
manager = self._manager((retained, removed, other, manual))
latencyScheduler = manager._latencyScheduler
pool = _ControlledThreadPool()
latencyScheduler.threadPool = pool
latencyScheduler.pingWorkerFactory = _ControlledLatencyWorker
manager.testPing((retained, removed, other))
manager.testDownloadSpeed(
(retained, removed, other),
concurrent=False,
)
manager.testDownloadSpeed((retained,), concurrent=True)
processQtEvents()
activeLatency = pool.started[0]
serialWorker = next(
worker
for worker in _ControlledDownloadWorker.instances
if worker.port in manager.SerialDownloadPorts
)
concurrentWorker = next(
worker
for worker in _ControlledDownloadWorker.instances
if worker.port in manager.ConcurrentDownloadPorts
)
self._setProfiles((retained, replacement, other, manual))
manager.invalidateSubscriptions({'group-a'}, clearResults=True)
self.assertIs(activeLatency.job.state, ProfileTestJobState.Cancelled)
self.assertEqual(serialWorker.cancelCount, 1)
self.assertEqual(concurrentWorker.cancelCount, 1)
self.assertTrue(
all(
job.target.subscriptionSource == 'group-b'
for job in latencyScheduler.queue
)
)
self.assertTrue(
all(
job.target.subscriptionSource == 'group-b'
for job in manager._serialDownloadScheduler.queue
)
)
self.assertFalse(manager._concurrentDownloadScheduler.queue)
self.assertEqual(retained.metadata.latency, '')
self.assertEqual(retained.metadata.speed, '')
self.assertEqual(replacement.metadata.latency, '')
self.assertEqual(replacement.metadata.speed, '')
self.assertEqual(other.metadata.latency, 'other latency')
self.assertEqual(other.metadata.speed, 'other speed')
self.assertEqual(manual.metadata.latency, 'manual latency')
self.assertEqual(manual.metadata.speed, 'manual speed')
activeLatency.finish('stale')
processQtEvents()
self.assertEqual(retained.metadata.latency, '')
self.assertEqual(len(pool.started), 2)
otherLatency = pool.started[1]
otherLatency.finish('9ms')
processQtEvents()
otherDownload = next(
worker
for worker in _ControlledDownloadWorker.instances
if worker.profile.metadata.profileId == other.metadata.profileId
)
otherDownload.finish('3.00 MiB/s')
processQtEvents()
self.assertEqual(other.metadata.latency, '9ms')
self.assertEqual(other.metadata.speed, '3.00 MiB/s')
self.assertFalse(
any(
worker.profile.metadata.profileId == replacement.metadata.profileId
for worker in _ControlledDownloadWorker.instances
)
)
def testWorkerReturnsResultsWithoutMutatingItsProfileSnapshot(self):
"""Keep worker execution independent from persistence write-back."""
profile = self._profile('profile', 'profile.example')
worker = _ControlledDownloadWorker(
profile.deepcopy(),
30000,
DownloadSpeedTestOptions(5000, 'https://example.test'),
)
results = []
worker.progressed.connect(lambda _worker, result: results.append(result))
worker.publish('2.00 MiB/s')
self.assertEqual(worker.profile.metadata.speed, '')
self.assertEqual(results[-1].value, '2.00 MiB/s')
worker.deleteLater()
def testShutdownDoesNotCommitAnActiveDownloadProgressValue(self):
"""Treat service shutdown as cancellation rather than result write-back."""
profile = self._profile('profile', 'profile.example')
manager = self._manager((profile,))
manager.testDownloadSpeed((profile,), concurrent=False)
processQtEvents()
worker = _ControlledDownloadWorker.instances[0]
worker.publish('1.00 MiB/s')
self.assertEqual(profile.metadata.speed, '1.00 MiB/s')
profile.metadata.speed = ''
manager.shutdown()
processQtEvents()
self.assertEqual(profile.metadata.speed, '')
def testRuntimeCallbacksAreDisposedBeforeWorkerDeferredDeletion(self):
"""Make a late core-exit callback harmless after terminal disposal."""
profile = self._profile('profile', 'profile.example')
worker = _DownloadSpeedWorker(
profile.deepcopy(),
30000,
DownloadSpeedTestOptions(5000, 'https://example.test'),
)
coreManager = _CoreManagerProbe()
worker.coreManager = coreManager
worker.finished.connect(lambda current, _result: current.deleteLater())
worker.runCompletionCallback()
self.assertEqual(coreManager.stopCount, 1)
self.assertTrue(waitFor(lambda: not isValid(worker)))
worker.coreExitCallback(profile.connection, 1)
self.assertEqual(profile.metadata.speed, '')
def testServerTableRepaintsOnlyTheCellCommittedByTheService(self):
"""Keep the one UI integration boundary limited to presentation."""
with isolatedSettings():
profile = self._profile('profile', 'profile.example')
Storage.UserServers().append(profile)
profile.index = 0
AppSettings.set('ActivatedItemIndex', '-1')
table = ServerTableView(
configurationEditorFactory=QWidget,
qrCodeWindowFactory=QWidget,
importActionsFactory=tuple,
)
try:
target = ProfileTestTarget.capture(profile)
result = ProfileTestResult(
ProfileTestField.DownloadSpeed,
'4.00 MiB/s',
)
with mock.patch.object(table, 'flushItem') as repaint:
self.assertTrue(
table.profileTestManager.applyResult(target, result)
)
self.assertEqual(profile.metadata.speed, '4.00 MiB/s')
repaint.assert_called_once_with(
0,
table.Headers.index('Speed'),
profile,
)
finally:
table.cleanup()
table.deleteLater()
Storage.UserServers().clear()
Storage._UserServersStorage.cache_clear()
Storage._UserSubsStorage.cache_clear()
processQtEvents()
if __name__ == '__main__':
unittest.main()
+11 -7
View File
@@ -17,24 +17,25 @@
"""Verify subscription workflow ownership outside table widgets."""
from types import SimpleNamespace
from unittest import TestCase, mock
import os
os.environ.setdefault('QT_QPA_PLATFORM', 'offscreen')
from PySide6 import QtCore, QtTest, QtWidgets
from shiboken6 import isValid
from Furious.Repository import Storage
from Furious.Repository.Subscriptions import SubscriptionGroup
from Furious.Service.SubscriptionManager import SubscriptionManager
from Furious.Window.SubscriptionPage import SubscriptionPage
from Furious.Widget.SubscriptionTableView import SubscriptionTableView
from tests.support import application, processQtEvents
from PySide6 import QtCore, QtTest, QtWidgets
from shiboken6 import isValid
from types import SimpleNamespace
from unittest import TestCase, mock
class _Payload:
"""Expose the minimal QNetworkReply byte-array contract."""
@@ -337,7 +338,9 @@ class SubscriptionManagerTest(TestCase):
side_effect=(RuntimeError('injected failure'), committed)
)
completed = []
committedSubscriptions = []
manager.updateCompleted.connect(completed.append)
manager.subscriptionCommitted.connect(committedSubscriptions.append)
failed = {'unique': 'group-a', 'profiles': ()}
successful = {'unique': 'group-b', 'profiles': ()}
@@ -355,6 +358,7 @@ class SubscriptionManagerTest(TestCase):
self.assertEqual(completed[0].successful[0]['unique'], 'group-b')
self.assertEqual(completed[0].failed[0]['unique'], 'group-a')
self.assertIn('injected failure', completed[0].failed[0]['error'])
self.assertEqual(committedSubscriptions, ['group-b'])
manager.deleteLater()