From 3844a8b9aae9e75d6941e40bb1ec1b44fb796285 Mon Sep 17 00:00:00 2001 From: Loren Eteval Date: Sun, 30 Aug 2026 12:46:43 +0800 Subject: [PATCH] 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 --- Furious/Core/CoreProcessWorker.py | 8 +- Furious/Service/AGENTS.md | 14 + Furious/Service/ProfileTesting.py | 1373 ++++++++++++++++++++++++ Furious/Service/SubscriptionManager.py | 6 + Furious/Service/TcpingService.py | 347 ++++++ Furious/Service/__init__.py | 18 + Furious/Widget/AGENTS.md | 5 +- Furious/Widget/ServerTableView.py | 761 ++----------- tests/README.md | 5 +- tests/test_profile_test_jobs.py | 926 ++++++++++++++++ tests/test_subscription_manager.py | 18 +- 11 files changed, 2783 insertions(+), 698 deletions(-) create mode 100644 Furious/Service/ProfileTesting.py create mode 100644 Furious/Service/TcpingService.py create mode 100644 tests/test_profile_test_jobs.py diff --git a/Furious/Core/CoreProcessWorker.py b/Furious/Core/CoreProcessWorker.py index c268e040..4460d001 100644 --- a/Furious/Core/CoreProcessWorker.py +++ b/Furious/Core/CoreProcessWorker.py @@ -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()), ) diff --git a/Furious/Service/AGENTS.md b/Furious/Service/AGENTS.md index db44655c..bd0ea7ba 100644 --- a/Furious/Service/AGENTS.md +++ b/Furious/Service/AGENTS.md @@ -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 diff --git a/Furious/Service/ProfileTesting.py b/Furious/Service/ProfileTesting.py new file mode 100644 index 00000000..238e5ba6 --- /dev/null +++ b/Furious/Service/ProfileTesting.py @@ -0,0 +1,1373 @@ +# Copyright (C) 2024–present Loren Eteval & contributors +# +# 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 . + +"""Provide identity-safe latency and download-speed profile testing.""" + +from __future__ import annotations + +from Furious.Frozenlib import ( + APP, + AppBuiltinRouting, + AppLogManager, + AppSettings, + NETWORK_SPEED_TEST_URL, + OS_CPU_COUNT, + classname, +) +from Furious.Core import CoreLaunchSpec +from Furious.Interface import CoreRuntime +from Furious.Models import ServerProfile, profileConnectionFingerprint +from Furious.Plugins import getPluginRegistry +from Furious.Qt.HttpGetManager import HttpGetManager +from Furious.Qt.Signals import connectWeakly, singleShotWeakly +from Furious.Repository import Storage +from Furious.Service.ConnectionManager import ConnectionManager +from Furious.Service.LogManager import CORE_LOG_CATEGORY +from Furious.Service.TcpingService import ( + TcpingCancelEvent, + TcpingEngine, + TcpingRequest, + TcpingRequestBatchEvent, + TcpingResultEvent, + TcpingThread, +) + +from PySide6 import QtCore +from PySide6.QtNetwork import QNetworkReply + +from dataclasses import dataclass, replace +from enum import Enum +from typing import Callable, Iterable + +import icmplib +import weakref +import collections + +__all__ = [ + 'DownloadSpeedTestOptions', + 'LatencyTestOptions', + 'LatencyTestType', + 'ProfileTestField', + 'ProfileTestJobState', + 'ProfileTestManager', + 'ProfileTestResult', + 'ProfileTestTarget', +] + + +def _appIsExiting() -> bool: + """Return whether the application is absent or shutting down.""" + app = APP() + + if app is None: + return True + + isExiting = getattr(app, 'isExiting', None) + + return isExiting() if callable(isExiting) else True + + +class ProfileTestField(Enum): + """Identify the persisted profile result field produced by a test.""" + + Latency = 'latency' + DownloadSpeed = 'speed' + + +class ProfileTestJobState(Enum): + """Describe one scheduler-owned test job lifecycle.""" + + Pending = 'pending' + Running = 'running' + Completed = 'completed' + Cancelled = 'cancelled' + + +class LatencyTestType(Enum): + """Select the execution model used for a latency test.""" + + Ping = 'ping' + Tcping = 'tcping' + + +@dataclass(frozen=True) +class LatencyTestOptions: + """Describe how one latency target should be tested.""" + + testType: LatencyTestType + timeoutMilliseconds: int = 2000 + + def __post_init__(self): + """Normalize and validate externally supplied latency options.""" + object.__setattr__(self, 'testType', LatencyTestType(self.testType)) + + if self.timeoutMilliseconds <= 0: + raise ValueError('latency timeout must be positive') + + +@dataclass(frozen=True) +class DownloadSpeedTestOptions: + """Describe one download test independently of its profile identity.""" + + timeoutMilliseconds: int + testUrl: str + logActionMessage: bool = False + + def __post_init__(self): + """Validate explicit download-test parameters.""" + if self.timeoutMilliseconds <= 0: + raise ValueError('download timeout must be positive') + + if not isinstance(self.testUrl, str): + raise TypeError('download test URL must be a string') + + +@dataclass(frozen=True) +class ProfileTestResult: + """Carry one presentation-compatible outcome from execution to write-back.""" + + field: ProfileTestField + value: str + terminal: bool = True + + def __post_init__(self): + """Normalize enum and stored display value.""" + object.__setattr__(self, 'field', ProfileTestField(self.field)) + object.__setattr__(self, 'value', str(self.value)) + + +@dataclass(frozen=True) +class ProfileTestTarget: + """Identify one immutable connection snapshot by repository identity.""" + + profileId: str + subscriptionSource: str + connectionFingerprint: str + snapshot: ServerProfile + + @classmethod + def capture(cls, profile: ServerProfile): + """Capture exactly what a test should execute at enqueue time.""" + return cls( + profile.metadata.profileId, + profile.metadata.subscriptionSource, + profileConnectionFingerprint(profile), + profile.deepcopy(), + ) + + +def _currentTargets(profiles: Iterable[ServerProfile]): + """Index current valid profiles and their connection fingerprints.""" + result = {} + + for profile in profiles: + if profile.deleted: + continue + + try: + fingerprint = profileConnectionFingerprint(profile) + except (TypeError, ValueError): + continue + + result[profile.metadata.profileId] = (profile, fingerprint) + + return result + + +def _resolveTarget(target: ProfileTestTarget, currentTargets): + """Resolve a target only while the same logical connection still exists.""" + current = currentTargets.get(target.profileId) + + if current is None or current[1] != target.connectionFingerprint: + return None + + return current[0] + + +@dataclass +class _LatencyTestJob: + """Pair one stable target with explicit latency-test options.""" + + target: ProfileTestTarget + options: LatencyTestOptions + state: ProfileTestJobState = ProfileTestJobState.Pending + + +@dataclass +class _TcpingRequestGroup: + """Associate one endpoint probe with every target awaiting its result.""" + + request: TcpingRequest + endpointKey: tuple + jobs: collections.deque + result: object = None + networkCompleted: bool = False + + +PING_RESULT_EVENT_TYPE = QtCore.QEvent.Type(QtCore.QEvent.registerEventType()) + + +class _PingResultEvent(QtCore.QEvent): + """Carry one blocking Ping result back after QRunnable execution returns.""" + + def __init__(self, job, result): + """Initialize one immutable scheduler delivery.""" + super().__init__(PING_RESULT_EVENT_TYPE) + + self.job = job + self.result = result + + +class _PingWorker(QtCore.QRunnable): + """Measure one ICMP endpoint using an isolated profile snapshot.""" + + def __init__(self, job: _LatencyTestJob, resultReceiver): + """Retain only immutable work and a weak scheduler receiver.""" + super().__init__() + + self.job = job + self.resultReceiver = weakref.ref(resultReceiver) + + def run(self): + """Execute the bounded blocking ping outside the GUI thread.""" + profile = self.job.target.snapshot + timeoutSeconds = self.job.options.timeoutMilliseconds / 1000 + + try: + response = icmplib.ping( + profile.itemAddress, + count=1, + timeout=timeoutSeconds, + interval=1, + ) + except Exception as ex: + latency = classname(ex) + else: + if response.address and response.is_alive: + latency = f'{round(response.avg_rtt)}ms' + elif response.packet_loss == 1: + latency = 'Timeout' + else: + latency = 'Error' + + receiver = self.resultReceiver() + + if receiver is not None and not _appIsExiting(): + try: + QtCore.QCoreApplication.postEvent( + receiver, _PingResultEvent(self.job, latency) + ) + except RuntimeError: + pass + + +class _LatencyScheduler(QtCore.QObject): + """Own bounded Ping jobs and one coalescing TCPing networking thread.""" + + MaximumTcpingConcurrency = 8 + TcpingResultBatchSize = 64 + + def __init__( + self, + resolveTarget: Callable[[ProfileTestTarget], ServerProfile | None], + publishResult: Callable[[ProfileTestTarget, ProfileTestResult], bool], + *, + pingConcurrency: int, + tcpingConcurrency: int, + parent=None, + threadPool=None, + pingWorkerFactory=_PingWorker, + ): + """Initialize explicit dependencies and owned execution resources.""" + super().__init__(parent) + + self._resolveTarget = resolveTarget + self._publishResult = publishResult + self.maxConcurrency = max(int(pingConcurrency), 1) + self.tcpingMaxConcurrency = max(int(tcpingConcurrency), 1) + self.threadPool = threadPool + self._ownsThreadPool = self.threadPool is None + + if self.threadPool is None: + self.threadPool = QtCore.QThreadPool(self) + self.threadPool.setMaxThreadCount(self.maxConcurrency) + + self.pingWorkerFactory = pingWorkerFactory + self.queue = collections.deque() + self.activeJobs = {} + self.tcpingRequests = {} + self.tcpingEndpointRequests = {} + self.tcpingCompletionQueue = collections.deque() + self.nextTcpingRequestId = 1 + self.tcpingThread = None + self.tcpingEngine = None + self.drainScheduled = False + self.tcpingCompletionScheduled = False + self.shuttingDown = False + + def event(self, event): + """Apply network-thread TCPing results in this object's Qt thread.""" + if isinstance(event, _PingResultEvent): + self.handleWorkerFinished(event.job, event.result) + + return True + + if isinstance(event, TcpingResultEvent): + self.handleTcpingResult(event.requestId, event.result) + + return True + + return super().event(event) + + def enqueue(self, profiles, options: LatencyTestOptions): + """Capture profiles and route them to the selected execution model.""" + if options.testType is LatencyTestType.Tcping: + self.enqueueTcping(profiles, options) + + return + + self.queue.extend( + _LatencyTestJob(ProfileTestTarget.capture(profile), options) + for profile in profiles + ) + self.scheduleDrain() + + def ensureTcpingEngine(self): + """Lazily start the one networking event loop owned by this scheduler.""" + if self.tcpingEngine is not None: + return self.tcpingEngine + + if self.shuttingDown: + raise RuntimeError('cannot start TCPing after scheduler shutdown') + + engine = TcpingEngine(self, self.tcpingMaxConcurrency) + thread = TcpingThread(engine, self) + engine.moveToThread(thread) + + self.tcpingThread = thread + self.tcpingEngine = engine + + thread.start() + + return engine + + @staticmethod + def tcpingEndpoint(profile): + """Return one normalized and validated endpoint tuple.""" + address, port = ( + str(profile.itemAddress).strip(), + int(profile.itemPort.split(',')[0]), + ) + + if not address or not 1 <= port <= 65535: + raise ValueError('invalid TCP endpoint') + + return (address.casefold(), port), address, port + + def enqueueTcping(self, profiles, options: LatencyTestOptions): + """Coalesce identical endpoints into one network-thread request.""" + requests = [] + + for profile in profiles: + job = _LatencyTestJob(ProfileTestTarget.capture(profile), options) + + try: + endpointKey, address, port = self.tcpingEndpoint(job.target.snapshot) + except Exception as ex: + self.completeJob(job, classname(ex)) + + continue + + requestKey = (*endpointKey, options.timeoutMilliseconds) + group = self.tcpingEndpointRequests.get(requestKey) + + if group is None: + request = TcpingRequest( + self.nextTcpingRequestId, + address, + port, + options.timeoutMilliseconds, + ) + + self.nextTcpingRequestId += 1 + + group = _TcpingRequestGroup( + request, + requestKey, + collections.deque(), + ) + + self.tcpingRequests[request.requestId] = group + self.tcpingEndpointRequests[requestKey] = group + + requests.append(request) + + job.state = ProfileTestJobState.Running + group.jobs.append(job) + + if requests: + QtCore.QCoreApplication.postEvent( + self.ensureTcpingEngine(), TcpingRequestBatchEvent(requests) + ) + + def completeJob(self, job, value): + """Publish one terminal latency result through the manager boundary.""" + result = ProfileTestResult(ProfileTestField.Latency, value) + + if job.state is not ProfileTestJobState.Cancelled and self._publishResult( + job.target, result + ): + job.state = ProfileTestJobState.Completed + else: + job.state = ProfileTestJobState.Cancelled + + def handleTcpingResult(self, requestId: int, result): + """Queue one endpoint result for bounded result fan-out.""" + group = self.tcpingRequests.get(requestId) + + if group is None or group.networkCompleted: + return + + group.result = result + group.networkCompleted = True + + if self.tcpingEndpointRequests.get(group.endpointKey) is group: + self.tcpingEndpointRequests.pop(group.endpointKey, None) + + self.tcpingCompletionQueue.append(requestId) + self.scheduleTcpingCompletion() + + def scheduleTcpingCompletion(self): + """Schedule one bounded GUI-thread completion batch.""" + if self.tcpingCompletionScheduled: + return + + self.tcpingCompletionScheduled = True + + singleShotWeakly(0, self, 'drainTcpingResults') + + def drainTcpingResults(self): + """Apply at most one fixed result batch before yielding to Qt.""" + self.tcpingCompletionScheduled = False + + remaining = self.TcpingResultBatchSize + + while self.tcpingCompletionQueue and remaining: + requestId = self.tcpingCompletionQueue[0] + group = self.tcpingRequests.get(requestId) + + if group is None: + self.tcpingCompletionQueue.popleft() + + continue + + while group.jobs and remaining: + self.completeJob(group.jobs.popleft(), group.result) + + remaining -= 1 + + if group.jobs: + break + + self.tcpingRequests.pop(requestId, None) + self.tcpingCompletionQueue.popleft() + + if self.tcpingCompletionQueue: + self.scheduleTcpingCompletion() + + def scheduleDrain(self): + """Schedule one bounded GUI-thread Ping queue drain.""" + if self.drainScheduled: + return + + self.drainScheduled = True + + singleShotWeakly(0, self, 'drain') + + def drain(self): + """Start valid blocking Ping jobs within the private-pool limit.""" + self.drainScheduled = False + + if _appIsExiting(): + self.cancelAll() + + return + + while self.queue and len(self.activeJobs) < self.maxConcurrency: + job = self.queue.popleft() + + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + + continue + + job.state = ProfileTestJobState.Running + worker = self.pingWorkerFactory(job, self) + + self.activeJobs[id(job)] = (job, worker) + self.threadPool.start(worker) + + @QtCore.Slot(object, object) + def handleWorkerFinished(self, job, result): + """Release one worker and publish only a still-current result.""" + active = self.activeJobs.pop(id(job), None) + + if active is None: + return + + self.completeJob(job, result) + self.scheduleDrain() + + def reconcileProfiles(self): + """Invalidate every job whose stable target no longer resolves.""" + retained = collections.deque() + + for job in self.queue: + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + else: + retained.append(job) + + self.queue = retained + + for job, worker in list(self.activeJobs.values()): + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + cancel = getattr(worker, 'cancel', None) + + if callable(cancel): + cancel() + + self.discardTcpingJobs(lambda job: self._resolveTarget(job.target) is None) + self.scheduleDrain() + + def invalidateSubscriptions(self, subscriptionIds): + """Discard work captured from the specified subscriptions.""" + subscriptionIds = {str(value) for value in subscriptionIds if value} + retained = collections.deque() + + for job in self.queue: + if job.target.subscriptionSource in subscriptionIds: + job.state = ProfileTestJobState.Cancelled + else: + retained.append(job) + + self.queue = retained + + for job, worker in list(self.activeJobs.values()): + if job.target.subscriptionSource in subscriptionIds: + job.state = ProfileTestJobState.Cancelled + cancel = getattr(worker, 'cancel', None) + + if callable(cancel): + cancel() + + self.discardTcpingJobs( + lambda job: job.target.subscriptionSource in subscriptionIds + ) + self.scheduleDrain() + + def discardTcpingJobs(self, predicate): + """Cancel matching group members and abort empty endpoint requests.""" + cancelRequestIds = [] + + for requestId, group in list(self.tcpingRequests.items()): + retained = collections.deque() + + for job in group.jobs: + if predicate(job): + job.state = ProfileTestJobState.Cancelled + else: + retained.append(job) + + group.jobs = retained + + if not retained: + self.tcpingRequests.pop(requestId, None) + + if self.tcpingEndpointRequests.get(group.endpointKey) is group: + self.tcpingEndpointRequests.pop(group.endpointKey, None) + + if not group.networkCompleted: + cancelRequestIds.append(requestId) + + if cancelRequestIds and self.tcpingEngine is not None: + QtCore.QCoreApplication.postEvent( + self.tcpingEngine, TcpingCancelEvent(cancelRequestIds) + ) + + def cancelAll(self): + """Cancel pending work, active TCPing, and stale-mark active Ping.""" + for job in self.queue: + job.state = ProfileTestJobState.Cancelled + + self.queue.clear() + + for job, worker in list(self.activeJobs.values()): + job.state = ProfileTestJobState.Cancelled + cancel = getattr(worker, 'cancel', None) + + if callable(cancel): + cancel() + + self.discardTcpingJobs(lambda _job: True) + + def shutdown(self): + """Stop the dedicated TCPing event loop exactly once.""" + if self.shuttingDown: + return + + self.shuttingDown = True + self.cancelAll() + + if self._ownsThreadPool and not self.threadPool.waitForDone(5000): + raise RuntimeError('Ping worker pool did not stop') + + thread = self.tcpingThread + engine = self.tcpingEngine + + if thread is None or engine is None: + return + + if thread.isRunning(): + if not QtCore.QMetaObject.invokeMethod( + engine, 'shutdown', QtCore.Qt.ConnectionType.BlockingQueuedConnection + ): + raise RuntimeError('failed to invoke TCPing engine shutdown') + + thread.quit() + + if not thread.wait(3000): + raise RuntimeError('TCPing networking thread did not stop') + + self.tcpingEngine = None + self.tcpingThread = None + + thread.deleteLater() + + +class _DownloadSpeedWorker(HttpGetManager): + """Own one temporary core and proxied HTTP download test.""" + + CoreStartupGraceMilliseconds = CoreLaunchSpec.DefaultWaitTime + + progressed = QtCore.Signal(object, object) + finished = QtCore.Signal(object, object) + + def __init__( + self, + profile: ServerProfile, + port: int, + options: DownloadSpeedTestOptions, + parent=None, + ): + """Initialize explicit test inputs and transient Qt resources.""" + super().__init__(parent, actionMessage='test download speed') + + self.profile = profile + self.port = port + self.options = options + self.result = ProfileTestResult( + ProfileTestField.DownloadSpeed, + '', + terminal=False, + ) + self.hasSpeedResult = False + self.totalBytesRead = 0 + self.hasDataCounter = 0 + self.cancelled = False + self._startInProgress = False + self._completionInProgress = False + self._pendingCompletionKwargs = None + self.coreManager = ConnectionManager() + self.networkReply = None + + self.elapsedTimer = QtCore.QElapsedTimer() + + self.coreStartupTimer = QtCore.QTimer(self) + self.coreStartupTimer.setSingleShot(True) + + self.timeoutTimer = QtCore.QTimer(self) + self.timeoutTimer.setSingleShot(True) + + connectWeakly( + self.coreStartupTimer.timeout, + self, + 'startDownload', + sender=self.coreStartupTimer, + ) + connectWeakly( + self.timeoutTimer.timeout, + self, + 'handleTimeout', + sender=self.timeoutTimer, + ) + + def setResult(self, value, *, publish=True): + """Record one semantic outcome without mutating a profile object.""" + self.result = ProfileTestResult( + ProfileTestField.DownloadSpeed, value, terminal=False + ) + + if publish: + self.publishProgress() + + def publishProgress(self): + """Publish a non-terminal result while this worker remains current.""" + if not self.cancelled and not self.completionHasRun and not _appIsExiting(): + self.progressed.emit(self, self.result) + + def completionCallback(self, **_kwargs): + """Dispose runtime callbacks before publishing terminal completion.""" + self.coreStartupTimer.stop() + self.timeoutTimer.stop() + + try: + self.coreManager.stopAll() + finally: + self.finished.emit(self, replace(self.result, terminal=True)) + + def runCompletionCallback(self, **kwargs): + """Defer terminal publication until synchronous startup has unwound.""" + if self.completionHasRun or self._completionInProgress: + return + + if self._startInProgress: + if self._pendingCompletionKwargs is None: + self._pendingCompletionKwargs = dict(kwargs) + + return + + self._completionInProgress = True + + try: + super().runCompletionCallback(**kwargs) + finally: + self._completionInProgress = False + + def _completionRequested(self) -> bool: + """Return whether startup must stop before acquiring another resource.""" + return ( + self.cancelled + or self.completionHasRun + or self._pendingCompletionKwargs is not None + ) + + def _finishStartPhase(self, phaseContinues: bool): + """Unwind one startup phase before honoring deferred completion.""" + pendingCompletionKwargs = self._pendingCompletionKwargs + + self._startInProgress = False + self._pendingCompletionKwargs = None + + if not phaseContinues or pendingCompletionKwargs is not None: + self.runCompletionCallback(**(pendingCompletionKwargs or {})) + + def isFinished(self) -> bool: + """Return whether the HTTP operation has no active reply.""" + if isinstance(self.networkReply, QNetworkReply): + return self.networkReply.isFinished() + + return True + + def abort(self): + """Abort the exact active HTTP reply if one exists.""" + if isinstance(self.networkReply, QNetworkReply): + self.networkReply.abort() + + def cancel(self): + """Cancel network/core work and complete exactly once.""" + if self.completionHasRun: + return + + self.cancelled = True + self.coreStartupTimer.stop() + self.timeoutTimer.stop() + + if not self.isFinished(): + self.abort() + + self.runCompletionCallback() + + def handleTimeout(self): + """Abort work that exceeds its explicit deadline.""" + try: + if not self.isFinished(): + self.abort() + finally: + self.runCompletionCallback() + + def coreExitCallback(self, _config, exitcode: int): + """Translate an unexpected temporary-core exit into a result.""" + if self.cancelled or self.completionHasRun: + return + + try: + if exitcode == CoreRuntime.ExitCode.ConfigurationError.value: + self.setResult('Invalid') + elif exitcode == CoreRuntime.ExitCode.ServerStartFailure.value: + self.setResult('Core start failed') + elif exitcode != CoreRuntime.ExitCode.SystemShuttingDown.value: + self.setResult(f'Core exited {exitcode}') + finally: + self.runCompletionCallback() + + def _startCoreRuntime(self) -> bool: + """Prepare and start one proxy-only temporary core.""" + config = getPluginRegistry().prepareDownloadTest(self.profile, self.port) + + if config is None: + self.setResult('Invalid') + + return False + + self.setResult('Starting') + + return self.coreManager.start( + config, + AppBuiltinRouting.Global.value, + self.coreExitCallback, + msgCallbackCore=AppLogManager().callback(CORE_LOG_CATEGORY), + deepcopy=False, + proxyModeOnly=True, + log=False, + waitCore=False, + ) + + def start(self): + """Launch the temporary core without blocking concurrent admission.""" + if self.completionHasRun or self._startInProgress: + return + + self._startInProgress = True + + readinessScheduled = False + + try: + if _appIsExiting() or self._completionRequested(): + return + + if not self.profile.isValid(): + self.setResult('Invalid') + elif self._startCoreRuntime() and not ( + _appIsExiting() or self._completionRequested() + ): + self.coreStartupTimer.start(self.CoreStartupGraceMilliseconds) + + readinessScheduled = True + finally: + self._finishStartPhase(readinessScheduled) + + @QtCore.Slot() + def startDownload(self): + """Start HTTP only after the concurrently launched core can become ready.""" + if self.completionHasRun or self._startInProgress: + return + + self._startInProgress = True + + downloadStarted = False + + try: + if _appIsExiting() or self._completionRequested(): + return + + if not self.coreManager.allRunning(): + self.setResult('Core start failed') + + return + + self.configureHttpProxy(f'127.0.0.1:{self.port}') + + if self._completionRequested(): + return + + self.networkReply = self.webGET( + self.options.testUrl, + logActionMessage=self.options.logActionMessage, + ) + + if self._completionRequested(): + if not self.isFinished(): + self.abort() + + return + + self.elapsedTimer.start() + self.timeoutTimer.start(self.options.timeoutMilliseconds) + + downloadStarted = True + finally: + self._finishStartPhase(downloadStarted) + + def _currentSpeed(self): + """Return a presentation-compatible MiB/s result.""" + elapsedSeconds = self.elapsedTimer.elapsed() / 1000 + speed = self.totalBytesRead / elapsedSeconds / 1024 / 1024 + + return f'{speed:.2f} MiB/s' + + def successCallback(self, networkReply, **_kwargs): + """Publish the final download speed or core-start failure.""" + if self.cancelled or self.completionHasRun: + return + + if self.coreManager.allRunning(): + self.totalBytesRead += networkReply.readAll().length() + self.setResult(self._currentSpeed(), publish=False) + else: + self.setResult('Core start failed', publish=False) + + self.coreManager.stopAll() + self.publishProgress() + + def hasDataCallback(self, networkReply, **_kwargs): + """Update bounded intermediate progress from newly available data.""" + if self.cancelled or self.completionHasRun: + return + + self.hasDataCounter += 1 + + if self.coreManager.allRunning(): + self.totalBytesRead += networkReply.readAll().length() + self.hasSpeedResult = True + self.setResult(self._currentSpeed(), publish=False) + + if self.hasDataCounter % 25 == 0: + self.publishProgress() + + def failureCallback(self, networkReply, **_kwargs): + """Translate one HTTP failure without exposing Qt enum details.""" + if self.cancelled or self.completionHasRun: + return + + if not self.hasSpeedResult: + if not self.coreManager.allRunning(): + return + + if ( + networkReply.error() + == QNetworkReply.NetworkError.OperationCanceledError + ): + value = 'Canceled' + else: + try: + value = networkReply.error().name + except Exception: + value = 'UnknownError' + + if isinstance(value, bytes): + value = value.decode('utf-8', 'replace') + elif not isinstance(value, str): + value = 'UnknownError' + + if value != 'UnknownError' and value.endswith('Error'): + value = value[:-5] + + self.setResult(value, publish=False) + + self.coreManager.stopAll() + self.publishProgress() + + +@dataclass +class _DownloadSpeedTestJob: + """Pair one stable target with explicit download-test options.""" + + target: ProfileTestTarget + options: DownloadSpeedTestOptions + state: ProfileTestJobState = ProfileTestJobState.Pending + + +class _DownloadSpeedScheduler(QtCore.QObject): + """Schedule one serial or concurrent stream of download jobs.""" + + def __init__( + self, + resolveTarget: Callable[[ProfileTestTarget], ServerProfile | None], + publishResult: Callable[[ProfileTestTarget, ProfileTestResult], bool], + *, + maxConcurrency: int, + portRange: range, + parent=None, + workerFactory=_DownloadSpeedWorker, + ): + """Initialize explicit identity, result, concurrency, and port inputs.""" + super().__init__(parent) + + if not isinstance(portRange, range) or len(portRange) == 0: + raise ValueError('download test port range cannot be empty') + + self._resolveTarget = resolveTarget + self._publishResult = publishResult + self.maxConcurrency = max(int(maxConcurrency), 1) + self.portRange = portRange + self.workerFactory = workerFactory + self.queue = collections.deque() + self.activeJobs = {} + self.activePorts = set() + self.nextPort = portRange.start + self.drainScheduled = False + + def enqueue(self, profiles, options: DownloadSpeedTestOptions): + """Capture each profile with the same explicit operation options.""" + self.queue.extend( + _DownloadSpeedTestJob(ProfileTestTarget.capture(profile), options) + for profile in profiles + ) + self.scheduleDrain() + + def cancelAll(self): + """Cancel every pending and active job through one terminal path.""" + for job in self.queue: + job.state = ProfileTestJobState.Cancelled + + self.queue.clear() + + for worker, job, _ in list(self.activeJobs.values()): + job.state = ProfileTestJobState.Cancelled + + worker.cancel() + + def scheduleDrain(self): + """Schedule valid pending jobs without recursive startup.""" + if self.drainScheduled: + return + + self.drainScheduled = True + + singleShotWeakly(0, self, 'drain') + + def drain(self): + """Start valid jobs while concurrency and local ports are available.""" + self.drainScheduled = False + + if _appIsExiting(): + self.cancelAll() + + return + + while self.queue and len(self.activeJobs) < self.maxConcurrency: + job = self.queue.popleft() + + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + + continue + + port = self.allocatePort() + + if port is None: + self.queue.appendleft(job) + + break + + self.startJob(job, port) + + def allocatePort(self): + """Reserve the next free port from this scheduler's explicit range.""" + for _ in range(len(self.portRange)): + port = self.nextPort + + self.nextPort += 1 + + if self.nextPort >= self.portRange.stop: + self.nextPort = self.portRange.start + + if port not in self.activePorts: + self.activePorts.add(port) + + return port + + return None + + def startJob(self, job: _DownloadSpeedTestJob, port: int): + """Construct one scheduler-owned transient worker and start it.""" + job.state = ProfileTestJobState.Running + + worker = self.workerFactory( + job.target.snapshot, + port, + job.options, + parent=self, + ) + + self.activeJobs[id(worker)] = (worker, job, port) + + connectWeakly( + worker.progressed, + self, + 'handleWorkerProgressed', + sender=worker, + ) + connectWeakly( + worker.finished, + self, + 'handleWorkerFinished', + sender=worker, + ) + + worker.start() + + @QtCore.Slot(object, object) + def handleWorkerProgressed(self, worker, result): + """Publish progress only while the exact target remains current.""" + active = self.activeJobs.get(id(worker)) + + if active is None: + return + + _, job, _ = active + + if not self._publishResult(job.target, result): + job.state = ProfileTestJobState.Cancelled + + worker.cancel() + + @QtCore.Slot(object, object) + def handleWorkerFinished(self, worker, result): + """Release one terminal worker after disposing its runtime callbacks.""" + active = self.activeJobs.pop(id(worker), None) + + if active is None: + return + + _, job, port = active + + if job.state is not ProfileTestJobState.Cancelled: + if self._publishResult(job.target, result): + job.state = ProfileTestJobState.Completed + else: + job.state = ProfileTestJobState.Cancelled + + self.activePorts.discard(port) + worker.deleteLater() + self.scheduleDrain() + + def reconcileProfiles(self): + """Remove pending and active jobs whose identities became stale.""" + retained = collections.deque() + + for job in self.queue: + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + else: + retained.append(job) + + self.queue = retained + + for worker, job, _port in list(self.activeJobs.values()): + if self._resolveTarget(job.target) is None: + job.state = ProfileTestJobState.Cancelled + + worker.cancel() + + self.scheduleDrain() + + def invalidateSubscriptions(self, subscriptionIds): + """Remove queued and active work owned by selected subscriptions.""" + subscriptionIds = {str(value) for value in subscriptionIds if value} + retained = collections.deque() + + for job in self.queue: + if job.target.subscriptionSource in subscriptionIds: + job.state = ProfileTestJobState.Cancelled + else: + retained.append(job) + + self.queue = retained + + for worker, job, _port in list(self.activeJobs.values()): + if job.target.subscriptionSource in subscriptionIds: + job.state = ProfileTestJobState.Cancelled + + worker.cancel() + + self.scheduleDrain() + + +class ProfileTestManager(QtCore.QObject): + """Own profile-test identity, execution, cancellation, and write-back.""" + + resultApplied = QtCore.Signal(object, object) + + SerialDownloadPorts = range(20809, 20810) + ConcurrentDownloadPorts = range(30000, 40000) + + def __init__( + self, + parent=None, + *, + profilesProvider: Callable[[], Iterable[ServerProfile]] | None = None, + pingConcurrency: int | None = None, + tcpingConcurrency: int | None = None, + downloadConcurrency: int | None = None, + ): + """Construct long-lived schedulers beneath one service owner.""" + super().__init__(parent) + + halfCpu = max(OS_CPU_COUNT // 2, 1) + pingConcurrency = halfCpu if pingConcurrency is None else pingConcurrency + tcpingConcurrency = ( + min(halfCpu, _LatencyScheduler.MaximumTcpingConcurrency) + if tcpingConcurrency is None + else tcpingConcurrency + ) + downloadConcurrency = ( + halfCpu if downloadConcurrency is None else downloadConcurrency + ) + + self._profilesProvider = profilesProvider or Storage.UserServers + self._targets = _currentTargets(self._profilesProvider()) + self._latencyScheduler = _LatencyScheduler( + self.resolveTarget, + self.applyResult, + pingConcurrency=pingConcurrency, + tcpingConcurrency=tcpingConcurrency, + parent=self, + ) + self._serialDownloadScheduler = _DownloadSpeedScheduler( + self.resolveTarget, + self.applyResult, + maxConcurrency=1, + portRange=self.SerialDownloadPorts, + parent=self, + ) + self._concurrentDownloadScheduler = _DownloadSpeedScheduler( + self.resolveTarget, + self.applyResult, + maxConcurrency=downloadConcurrency, + portRange=self.ConcurrentDownloadPorts, + parent=self, + ) + self._shuttingDown = False + + def resolveTarget(self, target: ProfileTestTarget): + """Resolve a captured target through the mutation-refreshed identity map.""" + return _resolveTarget(target, self._targets) + + def applyResult(self, target: ProfileTestTarget, result: ProfileTestResult) -> bool: + """Validate and write one result at the subsystem's sole commit boundary.""" + profile = self.resolveTarget(target) + + if profile is None: + return False + + self._writeResult(profile, result) + + return True + + def _writeResult(self, profile: ServerProfile, result: ProfileTestResult): + """Persist one compatible display value and notify presentation owners.""" + setattr(profile.metadata, result.field.value, result.value) + + self.resultApplied.emit(profile, result) + + def testPing(self, profiles, *, timeoutMilliseconds=2000): + """Queue ICMP latency tests for immutable profile snapshots.""" + self._latencyScheduler.enqueue( + profiles, + LatencyTestOptions(LatencyTestType.Ping, timeoutMilliseconds), + ) + + def testTcping(self, profiles, *, timeoutMilliseconds=2000): + """Queue coalesced asynchronous TCP latency tests.""" + self._latencyScheduler.enqueue( + profiles, + LatencyTestOptions(LatencyTestType.Tcping, timeoutMilliseconds), + ) + + def testDownloadSpeed( + self, + profiles, + *, + timeoutMilliseconds=5000, + concurrent=True, + testUrl=None, + logActionMessage=False, + ): + """Queue serial or concurrent downloads with explicit operation options.""" + if testUrl is None: + try: + configuredUrl = AppSettings.get('CustomNetworkSpeedTestURL') + except AttributeError: + configuredUrl = None + + testUrl = ( + configuredUrl + if isinstance(configuredUrl, str) + else NETWORK_SPEED_TEST_URL + ) + + options = DownloadSpeedTestOptions( + timeoutMilliseconds, + testUrl, + bool(logActionMessage), + ) + + scheduler = ( + self._concurrentDownloadScheduler + if concurrent + else self._serialDownloadScheduler + ) + scheduler.enqueue(profiles, options) + + def clearResults(self, profiles): + """Clear both presentation-compatible result fields for current profiles.""" + results = ( + ProfileTestResult(ProfileTestField.Latency, ''), + ProfileTestResult(ProfileTestField.DownloadSpeed, ''), + ) + + for profile in profiles: + for result in results: + self._writeResult(profile, result) + + def reconcileProfiles(self): + """Refresh current identity once and proactively invalidate stale work.""" + self._targets = _currentTargets(self._profilesProvider()) + self._latencyScheduler.reconcileProfiles() + self._serialDownloadScheduler.reconcileProfiles() + self._concurrentDownloadScheduler.reconcileProfiles() + + def invalidateSubscriptions(self, subscriptionIds, *, clearResults=False): + """Invalidate selected groups without disturbing unrelated work.""" + subscriptionIds = {str(value) for value in subscriptionIds if value} + + self._latencyScheduler.invalidateSubscriptions(subscriptionIds) + self._serialDownloadScheduler.invalidateSubscriptions(subscriptionIds) + self._concurrentDownloadScheduler.invalidateSubscriptions(subscriptionIds) + + if clearResults: + profiles = [ + profile + for profile in self._profilesProvider() + if profile.metadata.subscriptionSource in subscriptionIds + ] + self.clearResults(profiles) + + self.reconcileProfiles() + + def shutdown(self): + """Cancel exact work and stop every owned reusable execution resource.""" + if self._shuttingDown: + return + + self._shuttingDown = True + self._latencyScheduler.shutdown() + self._serialDownloadScheduler.cancelAll() + self._concurrentDownloadScheduler.cancelAll() diff --git a/Furious/Service/SubscriptionManager.py b/Furious/Service/SubscriptionManager.py index ab14aee9..7862d4fc 100644 --- a/Furious/Service/SubscriptionManager.py +++ b/Furious/Service/SubscriptionManager.py @@ -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: diff --git a/Furious/Service/TcpingService.py b/Furious/Service/TcpingService.py new file mode 100644 index 00000000..d8b87d11 --- /dev/null +++ b/Furious/Service/TcpingService.py @@ -0,0 +1,347 @@ +# Copyright (C) 2024–present Loren Eteval & contributors +# +# 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 . + +"""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 diff --git a/Furious/Service/__init__.py b/Furious/Service/__init__.py index 46ef01ab..75ea5a25 100644 --- a/Furious/Service/__init__.py +++ b/Furious/Service/__init__.py @@ -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', diff --git a/Furious/Widget/AGENTS.md b/Furious/Widget/AGENTS.md index d15e261b..629ba1a3 100644 --- a/Furious/Widget/AGENTS.md +++ b/Furious/Widget/AGENTS.md @@ -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. diff --git a/Furious/Widget/ServerTableView.py b/Furious/Widget/ServerTableView.py index 42117034..353e0961 100644 --- a/Furious/Widget/ServerTableView.py +++ b/Furious/Widget/ServerTableView.py @@ -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.""" diff --git a/tests/README.md b/tests/README.md index dceae928..34a932cc 100644 --- a/tests/README.md +++ b/tests/README.md @@ -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 ``` diff --git a/tests/test_profile_test_jobs.py b/tests/test_profile_test_jobs.py new file mode 100644 index 00000000..1818a556 --- /dev/null +++ b/tests/test_profile_test_jobs.py @@ -0,0 +1,926 @@ +# Copyright (C) 2024–present Loren Eteval & contributors +# +# 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 . + +"""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() diff --git a/tests/test_subscription_manager.py b/tests/test_subscription_manager.py index b5fb88b1..1858064f 100644 --- a/tests/test_subscription_manager.py +++ b/tests/test_subscription_manager.py @@ -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()