From c35ea16a29e5513ab375cbf8dcb9c6bb96c8a620 Mon Sep 17 00:00:00 2001 From: Loren Eteval Date: Wed, 7 Oct 2026 13:18:30 +0800 Subject: [PATCH] Fix Qt service lifetime cleanup Signed-off-by: Loren Eteval --- Furious/Controllers/AGENTS.md | 3 + Furious/Controllers/ConnectionController.py | 2 +- Furious/Service/AGENTS.md | 7 + Furious/Service/ConnectivityManager.py | 27 +- Furious/Service/DnsResolver.py | 14 + Furious/Service/ProfileTesting.py | 21 +- Furious/Service/TrafficStatsManager.py | 22 ++ tests/README.md | 5 + tests/fixtures/editor_lifetime_probe.py | 170 +++++++++++- tests/test_profile_test_jobs.py | 173 +++++++++++- tests/test_service_runtime.py | 284 +++++++++++++++++++- 11 files changed, 721 insertions(+), 7 deletions(-) diff --git a/Furious/Controllers/AGENTS.md b/Furious/Controllers/AGENTS.md index fa74d67a..0eef6837 100644 --- a/Furious/Controllers/AGENTS.md +++ b/Furious/Controllers/AGENTS.md @@ -9,6 +9,9 @@ compatibility paths. - Controllers own process-lifetime shared state and transition policy. They coordinate injected repositories/services and publish structured Qt signals. New behavior delegates presentation and execution resources to their owners; existing host-setting prompts do not justify moving network replies, core processes, or pools into controllers. + A default QObject service created by a controller is its native child; an injected service keeps its supplied + owner. `ConnectionController` update-manager teardown and borrowed-owner cases in `tests/test_service_runtime.py` + verify that a retained invalid controller wrapper cannot keep its default network requests alive. - `ConnectionController` is the sole connection state machine. A GUI start remains `Connecting` while one generation-checked `ConnectionManager` transaction acquires readiness/TUN resources. The selected live profile is exposed during `Connecting`; successful runtime commit precedes System Proxy setup and `Connected`. Failure resets diff --git a/Furious/Controllers/ConnectionController.py b/Furious/Controllers/ConnectionController.py index ef5a962b..af03043c 100644 --- a/Furious/Controllers/ConnectionController.py +++ b/Furious/Controllers/ConnectionController.py @@ -103,7 +103,7 @@ class ConnectionController(QtCore.QObject): self._actionQueue = queue.Queue() self._coreManager = coreManager or ConnectionManager() - self._updatesManager = updatesManager or UpdateManager() + self._updatesManager = updatesManager or UpdateManager(self) self._state = ConnectionState.Disconnected self._activeProfile = None self._lastError = None diff --git a/Furious/Service/AGENTS.md b/Furious/Service/AGENTS.md index 798a7c2c..528b519d 100644 --- a/Furious/Service/AGENTS.md +++ b/Furious/Service/AGENTS.md @@ -13,6 +13,13 @@ for execution, and Qt for lifetime primitives. This scope owns multi-stage workf owner and explicit idempotent cleanup. Cancellation can suppress a result without stopping the underlying work; distinguish deadline-bounded teardown from cooperative drains, and retain resources until their users finish. Construct Qt services only after an application exists. +- Native owner destruction also ends Python execution ownership. At `destroyed`, the owner's wrapper is invalid + but its QObject children have not yet been deleted; a plain weak-reference callback may release Python state + and shut down still-valid child schedulers without calling the destroyed owner's Qt API. Statistics executors, + profile-test runtime leases and DNS-operation replies exercise this boundary in `test_service_runtime.py` and + `test_profile_test_jobs.py`. This final attempt does not guarantee release of a resource that refuses cleanup: + explicit shutdown must retain retry ownership before native deletion. Cancellation of a running provider remains + cooperative, and executor admission closure is distinct from actual worker termination. - Inject repositories/providers/clients/runtime factories where practical. Stage results, prove freshness, and commit through the owning repository/controller rather than creating a parallel authoritative collection. - Every async workflow defines supersession and one terminal publication path. Generation/version or exact target diff --git a/Furious/Service/ConnectivityManager.py b/Furious/Service/ConnectivityManager.py index 17ea1e75..f92e7935 100644 --- a/Furious/Service/ConnectivityManager.py +++ b/Furious/Service/ConnectivityManager.py @@ -21,10 +21,13 @@ from __future__ import annotations from Furious.Frozenlib import * from Furious.Qt.HttpGetManager import * +from Furious.Qt.Signals import connectWeakly from PySide6 import QtCore from PySide6.QtNetwork import * +from shiboken6 import isValid + import logging __all__ = ['ConnectivityManager'] @@ -60,9 +63,23 @@ class ConnectivityManager(Mixins.ConnectionAware, HttpGetManager): @QtCore.Slot() def _abortActiveReply(self): """Abort the one currently active connectivity request.""" - if isinstance(self._activeReply, QNetworkReply): + if isinstance(self._activeReply, QNetworkReply) and isValid(self._activeReply): self._activeReply.abort() + @QtCore.Slot() + def _activeReplyDestroyed(self): + """Release a probe whose native reply disappeared without finishing.""" + if self._activeReply is None or isValid(self._activeReply): + return + + # A completed reply may be deleted after another probe has started. + # Only the invalid active reply owns this timeout and admission slot. + self._activeReply = None + self.jobTimeoutTimer.stop() + + if self._testingEnabled: + self.jobArrangeTimer.start(self.recalculateJobInterval(jobStatus=False)) + def recalculateJobInterval(self, jobStatus: bool) -> int: """Return the recalculate job interval value used by the network connectivity manager.""" assert isinstance(jobStatus, bool) @@ -114,6 +131,14 @@ class ConnectivityManager(Mixins.ConnectionAware, HttpGetManager): url = NETWORK_CONNECTIVITY_TEST_URL self._activeReply = self.webGET(url) + + connectWeakly( + self._activeReply.destroyed, + self, + '_activeReplyDestroyed', + sender=self._activeReply, + ) + self.jobTimeoutTimer.start(ConnectivityManager.MIN_JOB_INTERVAL - 500) def stopTest(self): diff --git a/Furious/Service/DnsResolver.py b/Furious/Service/DnsResolver.py index 5ed77970..53fa9f6a 100644 --- a/Furious/Service/DnsResolver.py +++ b/Furious/Service/DnsResolver.py @@ -31,13 +31,23 @@ from shiboken6 import isValid from typing import Tuple +import weakref import logging +import functools __all__ = ['DnsResolutionOperation', 'DnsResolver'] logger = logging.getLogger(__name__) +def _cancelDestroyedDnsOperation(operationReference, *_args): + """Abort operation-owned requests before Qt deletes its observer timer.""" + operation = operationReference() + + if operation is not None: + operation.cancel() + + class DnsResolutionOperation(QtCore.QObject): """Observe one recursive DNS request without nesting the Qt event loop.""" @@ -59,6 +69,10 @@ class DnsResolutionOperation(QtCore.QObject): connectWeakly(self._timer.timeout, self, '_poll') + self.destroyed.connect( + functools.partial(_cancelDestroyedDnsOperation, weakref.ref(self)) + ) + def start(self): """Start the DNS request and its event-driven completion observer.""" if self._terminal or self._timer.isActive(): diff --git a/Furious/Service/ProfileTesting.py b/Furious/Service/ProfileTesting.py index 1033625e..e422da43 100644 --- a/Furious/Service/ProfileTesting.py +++ b/Furious/Service/ProfileTesting.py @@ -50,6 +50,8 @@ from Furious.Service.TcpingService import ( from PySide6 import QtCore from PySide6.QtNetwork import QNetworkReply +from shiboken6 import isValid + from dataclasses import dataclass, replace from enum import Enum from typing import Callable, Iterable @@ -58,6 +60,7 @@ import icmplib import logging import weakref import collections +import functools logger = logging.getLogger(__name__) @@ -823,14 +826,14 @@ class _DownloadSpeedWorker(HttpGetManager): def isFinished(self) -> bool: """Return whether the HTTP operation has no active reply.""" - if isinstance(self.networkReply, QNetworkReply): + if isinstance(self.networkReply, QNetworkReply) and isValid(self.networkReply): return self.networkReply.isFinished() return True def abort(self): """Abort the exact active HTTP reply if one exists.""" - if isinstance(self.networkReply, QNetworkReply): + if isinstance(self.networkReply, QNetworkReply) and isValid(self.networkReply): self.networkReply.abort() def cancel(self): @@ -1325,6 +1328,16 @@ class _DownloadSpeedScheduler(QtCore.QObject): self.scheduleDrain() +def _shutdownDestroyedProfileTestManager(managerReference, *_args): + """Release execution before Qt deletes the manager's child schedulers.""" + manager = managerReference() + + if manager is not None: + # destroyed invalidates the manager, but its children are still valid. + # shutdown uses those children and Python state, never the manager's Qt API. + manager.shutdown() + + class ProfileTestManager(QtCore.QObject): """Own profile-test identity, execution, cancellation, and write-back.""" @@ -1385,6 +1398,10 @@ class ProfileTestManager(QtCore.QObject): self._shuttingDown = False + self.destroyed.connect( + functools.partial(_shutdownDestroyedProfileTestManager, weakref.ref(self)) + ) + def resolveTarget(self, target: ProfileTestTarget): """Resolve a captured target through the mutation-refreshed identity map.""" return _resolveTarget(target, self._targets) diff --git a/Furious/Service/TrafficStatsManager.py b/Furious/Service/TrafficStatsManager.py index 85a794d4..08f9868e 100644 --- a/Furious/Service/TrafficStatsManager.py +++ b/Furious/Service/TrafficStatsManager.py @@ -200,6 +200,19 @@ def _forwardTrafficQueryResult(managerReference, generation, future): manager._futureCompleted(generation, future) +def _closeDestroyedStatisticsManager(managerReference, *_args): + """Close Python worker admission after Qt destroys the native manager.""" + manager = managerReference() + + if manager is not None: + # Native teardown has already invalidated this wrapper. Release only + # Python state here; the QObject tree stops and deletes the sample timer. + manager._generation += 1 + manager._connected = False + manager._monitor = None + manager._closeExecutor() + + class TrafficStatsManager( Mixins.ConnectionAware, Mixins.CleanupOnExit, @@ -238,6 +251,10 @@ class TrafficStatsManager( self._sampleReady.connect(self._consumeResult) + self.destroyed.connect( + functools.partial(_closeDestroyedStatisticsManager, weakref.ref(self)) + ) + def _resumeSampling(self): """Resume polling whenever a statistics monitor is available.""" if not self._collectionEnabled or self._monitor is None: @@ -378,10 +395,15 @@ class TrafficStatsManager( return if self._connected: + generation = self._generation + monitor = getPluginRegistry().trafficStatsMonitorForRuntimes( self._activeRuntimes() ) + if not isValid(self) or generation != self._generation: + return + self._activateMonitor(monitor) @QtCore.Slot() diff --git a/tests/README.md b/tests/README.md index 0e674ef8..96720276 100644 --- a/tests/README.md +++ b/tests/README.md @@ -338,6 +338,11 @@ through native destruction, with weak-wrapper and compiled callback-retention ch Pool notifications skip peers destroyed by earlier callbacks. DNS observers cover reentrant cancellation/destruction during request creation and timeout aborts; statistics publication covers native manager destruction between emitted signals. +Service teardown additionally covers connectivity replies destroyed without `finished`, DNS-operation owner +destruction, statistics executor/thread release with an invalid wrapper retained, and download-runtime lease +release before native child-worker deletion. Controller-owned update requests die with their controller; +injected services retain their existing owner. Source regressions also exercise reentrant metrics enablement, +download cancellation with a deleted reply, and disposal during an output-drain callback. Endpoint lookup covers service destruction/disablement during state and result notifications, without admitting the next request from the abandoned stage. diff --git a/tests/fixtures/editor_lifetime_probe.py b/tests/fixtures/editor_lifetime_probe.py index 94aefbf0..b0f687f3 100644 --- a/tests/fixtures/editor_lifetime_probe.py +++ b/tests/fixtures/editor_lifetime_probe.py @@ -50,10 +50,17 @@ from Furious.Service.EndpointInfoService import ( ) from Furious.Service.SubscriptionManager import SubscriptionManager from Furious.Service.ProfileTesting import _LatencyScheduler +from Furious.Service.ProfileTesting import ( + DownloadSpeedTestOptions, + ProfileTestManager, + ProfileTestJobState, + _DownloadSpeedWorker, +) +from Furious.Service.ConnectivityManager import ConnectivityManager from Furious.Service.DnsResolver import DnsResolutionOperation, DnsResolver from Furious.Service.TrafficStatsManager import TrafficStatsManager from Furious.Frozenlib import Mixins -from Furious.Plugins import TrafficCounters +from Furious.Plugins import TrafficCounters, TrafficStatsMonitor from Furious.Service.ConnectionManager import ( ConnectionManager, ConnectionStartOperation, @@ -88,6 +95,7 @@ import argparse import json import sys import weakref +import threading from pathlib import Path import tempfile @@ -836,6 +844,165 @@ class _RequestPayload: """Expose whether pending request context still owns plain operation data.""" +def runServiceTeardownProbe(iterations=100): + """Exercise reply loss, executor teardown, and child runtime release compiled.""" + controller = ConnectionController(coreManager=SimpleNamespace(runtimes=[])) + updater = controller._updatesManager + reply = _PendingReply(updater) + + with mock.patch.object(updater, 'get', lambda _request: reply): + updater.webGET('https://invalid.test', marker='fixture') + + deleteQObject(controller) + + assert not isValid(updater) and not isValid(reply) + assert not updater._replyContexts + + manager = ConnectivityManager() + references = [] + + try: + manager._testingEnabled = True + + with mock.patch( + 'Furious.Service.ConnectivityManager.AppSettings.get', return_value=None + ): + for _ in range(iterations): + reply = _PendingReply(manager) + references.append(weakref.ref(reply)) + + with mock.patch.object(manager, 'get', lambda _request: reply): + manager.startSingleTest() + + assert manager._activeReply is reply + + deleteQObject(reply) + + assert manager._activeReply is None + assert not manager._replyContexts + assert not manager.jobTimeoutTimer.isActive() + + del reply + + assert all(reference() is None for reference in references) + finally: + manager.stopTest() + deleteQObject(manager) + + resolver = DnsResolver() + + try: + with mock.patch('Furious.Service.DnsResolver.logger.error'): + for _ in range(iterations): + parent = QtCore.QObject() + operation = resolver.resolveAsync('example.test', parent=parent) + reply = _PendingReply(resolver) + + with mock.patch.object(resolver, 'get', lambda _request: reply): + operation.start() + + deleteQObject(parent) + + assert not isValid(operation) + assert reply.isFinished() + assert not resolver._replyContexts + + processQtEvents() + + assert not isValid(reply) + finally: + resolver.dispose() + + parent = QtCore.QObject() + stats = TrafficStatsManager(parent) + started, release = threading.Event(), threading.Event() + + def query(_target): + started.set() + release.wait(5) + + return TrafficCounters(1, 2) + + stats._activateMonitor(TrafficStatsMonitor(query=query, target=None)) + executor = stats._executor + + try: + assert started.wait(2) + threads = tuple(executor._threads) + + deleteQObject(parent) + + assert stats._executor is None and stats._future is None + assert stats._monitor is None and executor._shutdown + + release.set() + + for thread in threads: + thread.join(3) + + assert all(not thread.is_alive() for thread in threads) + finally: + release.set() + executor.shutdown(wait=True, cancel_futures=True) + + if isValid(parent): + deleteQObject(parent) + + class Lease: + def __init__(self): + self.resourceOwned = True + + def release(self): + self.resourceOwned = False + + return True + + runtimeReleases = 0 + + for _ in range(iterations): + parent = QtCore.QObject() + tests = ProfileTestManager(parent, profilesProvider=tuple) + workers, leases = [], [] + + for scheduler in ( + tests._serialDownloadScheduler, + tests._concurrentDownloadScheduler, + ): + worker = _DownloadSpeedWorker( + None, + scheduler.portRange.start, + DownloadSpeedTestOptions(5000, 'https://invalid.test'), + parent=scheduler, + ) + lease = Lease() + worker._runtimeLease = lease + job = SimpleNamespace(state=ProfileTestJobState.Running) + scheduler.activeJobs[id(worker)] = (worker, job, worker.port) + scheduler.activePorts.add(worker.port) + connectWeakly( + worker.finished, scheduler, 'handleWorkerFinished', sender=worker + ) + workers.append(worker) + leases.append(lease) + + deleteQObject(parent) + processQtEvents() + + assert all(not lease.resourceOwned for lease in leases) + assert all(not isValid(worker) for worker in workers) + assert not tests._serialDownloadScheduler.activePorts + assert not tests._concurrentDownloadScheduler.activePorts + runtimeReleases += len(leases) + + return { + 'connectivityReplyDestructions': iterations, + 'controllerOwnedUpdateRequestsReleased': 1, + 'dnsRequestsAbortedOnOwnerDestruction': iterations, + 'statisticsThreadsStopped': len(threads), + 'downloadLeaseReleases': runtimeReleases, + } + + def runNetworkProbe(iterations=100): """Verify native teardown releases both reply registries under compilation.""" application() @@ -1491,6 +1658,7 @@ def main(): print(json.dumps(runButtonOwnershipProbe(arguments.iterations), sort_keys=True)) print(json.dumps(runConfirmationProbe(arguments.iterations), sort_keys=True)) print(json.dumps(runNetworkProbe(arguments.iterations), sort_keys=True)) + print(json.dumps(runServiceTeardownProbe(arguments.iterations), sort_keys=True)) print(json.dumps(runInfrastructureProbe(arguments.iterations), sort_keys=True)) print( diff --git a/tests/test_profile_test_jobs.py b/tests/test_profile_test_jobs.py index 6fa4128f..2e7e42ca 100644 --- a/tests/test_profile_test_jobs.py +++ b/tests/test_profile_test_jobs.py @@ -43,7 +43,7 @@ from Furious.Widget.ServerTableView import ServerTableView from PySide6 import QtCore, QtNetwork from PySide6.QtWidgets import QWidget -from shiboken6 import isValid +from shiboken6 import isValid, delete as deleteQObject from tests.support import ( application, @@ -279,6 +279,17 @@ class _CancelDuringStartDownloadWorker(_DownloadSpeedWorker): return True +class _PendingDownloadReply(QtNetwork.QNetworkReply): + """Own a real Qt reply without opening a network connection.""" + + def abort(self): + self.setFinished(True) + self.finished.emit() + + def readData(self, _maximumLength): + return bytes() + + class ProfileTestServiceTest(unittest.TestCase): """Exercise the self-contained profile-test subsystem.""" @@ -405,6 +416,40 @@ class ProfileTestServiceTest(unittest.TestCase): processQtEvents() + def testDeletedDownloadReplyDoesNotPreventCancellationOrTimeoutCleanup(self): + """Retained invalid replies cannot strand a worker's runtime and port.""" + for boundary in ('cancel', 'timeout'): + with self.subTest(boundary=boundary): + profile = self._profile('profile', 'example.test') + + with self._runtimeDownloads((profile,), failure=None) as ( + manager, + workers, + runtimes, + ): + manager.testDownloadSpeed((profile,), concurrent=False) + processQtEvents() + worker, runtime = workers[0], runtimes[0] + reply = _PendingDownloadReply(worker) + + with mock.patch.object(worker, 'get', return_value=reply): + worker.startDownload() + + deleteQObject(reply) + + if boundary == 'cancel': + manager.cancelAll() + else: + worker.timeoutTimer.timeout.emit() + + processQtEvents() + + self.assertFalse(runtime.isRunning()) + self.assertFalse(runtime.resourceOwned) + self.assertFalse(isValid(worker)) + self.assertFalse(manager._serialDownloadScheduler.activePorts) + self.assertFalse(manager._serialDownloadScheduler.activeJobs) + def testDownloadWorkerRetainsFailedLeaseUntilReleaseSucceeds(self): """A failed stop or dispose must leave the worker's exact lease reachable.""" for failure in ('stop', 'dispose'): @@ -538,6 +583,132 @@ class ProfileTestServiceTest(unittest.TestCase): self.assertFalse(manager._serialDownloadScheduler._pendingReleases) self.assertFalse(manager._serialDownloadScheduler.activePorts) + def testNativeManagerDestructionReleasesDownloadRuntimesBeforeWorkers(self): + """Qt owner teardown must release active and retained execution leases.""" + for pendingRelease in (False, True): + with self.subTest(pendingRelease=pendingRelease): + for _ in range(20): + profile = self._profile('profile', 'example.test') + parent = QtCore.QObject() + manager = ProfileTestManager( + parent, + profilesProvider=lambda: (profile,), + downloadConcurrency=1, + ) + workers, runtimes, routers, destroyed = [], [], [], [] + results, errors = [], [] + + def workerFactory(*args, **kwargs): + worker = _DownloadSpeedWorker(*args, **kwargs) + worker.CoreStartupGraceMilliseconds = 60_000 + worker.destroyed.connect(lambda *_args: destroyed.append(True)) + workers.append(worker) + + return worker + + def createRuntime(*_args, **kwargs): + runtime = _CleanupRefusingRuntime( + None, exitCallback=kwargs['exitCallback'] + ) + runtimes.append(runtime) + + return PreparedRuntime(runtime) + + schedulers = ( + manager._serialDownloadScheduler, + manager._concurrentDownloadScheduler, + ) + + for scheduler in schedulers: + scheduler.workerFactory = workerFactory + + registry = SimpleNamespace( + prepareDownloadTest=lambda profile, _port: profile, + createCoreRuntime=createRuntime, + ) + manager.resultApplied.connect(lambda *_args: results.append(True)) + + try: + with ( + mock.patch( + 'Furious.Service.ProfileTesting.getPluginRegistry', + return_value=registry, + ), + mock.patch('Furious.Service.ProfileTesting.AppLogManager'), + mock.patch( + 'sys.excepthook', lambda *args: errors.append(args) + ), + ): + manager.testDownloadSpeed((profile,), concurrent=False) + manager.testDownloadSpeed((profile,), concurrent=True) + processQtEvents() + + self.assertEqual(len(workers), 2) + routers.extend( + worker._runtimeLease.router for worker in workers + ) + self.assertTrue( + all(runtime.isRunning() for runtime in runtimes) + ) + + if pendingRelease: + for runtime, worker in zip(runtimes, workers): + runtime.failure = 'stop' + + with self.assertLogs( + 'Furious.Service.RuntimeLease', level='ERROR' + ): + worker.runCompletionCallback() + + with self.assertLogs( + 'Furious.Service.RuntimeLease', level='ERROR' + ): + processQtEvents() + + self.assertTrue( + all( + scheduler._pendingReleases + for scheduler in schedulers + ) + ) + + for runtime in runtimes: + runtime.failure = None + + results.clear() + deleteQObject(parent) + processQtEvents() + + self.assertFalse(errors) + self.assertFalse(results) + self.assertTrue( + all(not runtime.isRunning() for runtime in runtimes) + ) + self.assertTrue( + all(not runtime.resourceOwned for runtime in runtimes) + ) + self.assertTrue(manager._shuttingDown) + self.assertTrue(all(not isValid(worker) for worker in workers)) + self.assertTrue(all(not isValid(router) for router in routers)) + self.assertEqual(len(destroyed), 2) + + for scheduler in schedulers: + self.assertFalse(scheduler.activeJobs) + self.assertFalse(scheduler.activePorts) + self.assertFalse(scheduler._pendingReleases) + finally: + for runtime in runtimes: + runtime.failure = None + runtime.dispose() + + if isValid(manager): + manager.shutdown() + + if isValid(parent): + deleteQObject(parent) + + processQtEvents() + def testCancelAllPreservesResultsRejectsLatePingAndAllowsNewTests(self): """Cancel active and queued work across all schedulers without shutting down.""" profiles = [self._profile(str(i), f'{i}.example') for i in range(3)] diff --git a/tests/test_service_runtime.py b/tests/test_service_runtime.py index fa34e8ec..80e82abc 100644 --- a/tests/test_service_runtime.py +++ b/tests/test_service_runtime.py @@ -34,6 +34,7 @@ from Furious.Service.DnsResolver import DnsResolver from Furious.Repository import Storage from Furious.Service.TrafficStatsManager import TrafficStatsManager from Furious.Service.UpdateManager import UpdateManager +from Furious.Controllers.ConnectionController import ConnectionController from PySide6 import QtCore from PySide6.QtNetwork import QNetworkReply @@ -122,6 +123,48 @@ class UpdateManagerTest(unittest.TestCase): manager.deleteLater() parent.deleteLater() + def testDefaultUpdateManagerDiesWithItsController(self): + """An invalid controller wrapper cannot keep its update request alive.""" + controller = ConnectionController(coreManager=SimpleNamespace(runtimes=[])) + manager = controller._updatesManager + reply = _ManagedReply(manager) + payload = _ResponseBody(b'fixture') + reference = weakref.ref(payload) + + try: + with patch.object(manager, 'get', lambda _request: reply): + manager.webGET('https://invalid.test', payload=payload) + + del payload + deleteQObject(controller) + + self.assertFalse(isValid(manager)) + self.assertFalse(isValid(reply)) + self.assertFalse(manager._replyContexts) + self.assertIsNone(reference()) + finally: + if isValid(manager): + deleteQObject(manager) + + if isValid(controller): + deleteQObject(controller) + + def testInjectedUpdateManagerKeepsItsExistingOwner(self): + """Controller composition must not steal a borrowed network manager.""" + owner = QtCore.QObject() + manager = UpdateManager(owner) + controller = ConnectionController( + coreManager=SimpleNamespace(runtimes=[]), updatesManager=manager + ) + + try: + deleteQObject(controller) + + self.assertTrue(isValid(manager)) + self.assertIs(manager.parent(), owner) + finally: + deleteQObject(owner) + class _ManagedReply(QNetworkReply): """Provide a hermetic reply object with real Qt lifecycle signals.""" @@ -153,6 +196,54 @@ class _CapturingHttpGetManager(HttpGetManager): class HttpGetManagerLifetimeTest(unittest.TestCase): """Verify every request receives a timeout and releases exact context.""" + def testDnsOperationOwnerDestructionAbortsItsRequest(self): + """A resolver outliving a request must not outlive that request's owner.""" + application() + resolver = DnsResolver() + self.addCleanup(resolver.dispose) + + for _ in range(30): + parent = QtCore.QObject() + reply = _ManagedReply(resolver) + operation = resolver.resolveAsync('example.test', parent=parent) + results, errors = [], [] + operation.finished.connect(lambda *_args: results.append(True)) + + def abort(): + reply.setError( + QNetworkReply.NetworkError.OperationCanceledError, 'cancelled' + ) + reply.setFinished(True) + reply.finished.emit() + + reply.abort = abort + + try: + with ( + patch.object(resolver, 'get', lambda _request: reply), + patch('Furious.Service.DnsResolver.logger.error'), + patch('sys.excepthook', lambda *args: errors.append(args)), + ): + operation.start() + deleteQObject(parent) + + self.assertFalse(isValid(operation)) + self.assertFalse(isValid(operation._timer)) + self.assertTrue(reply.isFinished()) + self.assertFalse(resolver._replyContexts) + self.assertFalse(results) + self.assertFalse(errors) + + processQtEvents() + + self.assertFalse(isValid(reply)) + finally: + if isValid(reply): + reply.abort() + + if isValid(parent): + deleteQObject(parent) + def testDnsDisposalToleratesNativeDestructionDuringAbort(self): """Reply abort may destroy the resolver before its deferred delete.""" application() @@ -566,10 +657,97 @@ class ConnectivityManagerTest(unittest.TestCase): """Create the process-wide headless QApplication.""" application() + def testDestroyedActiveReplyReleasesProbeAndAllowsAnotherRequest(self): + """Native deletion without finished must not strand the probe scheduler.""" + manager = ConnectivityManager() + self.addCleanup(manager.deleteLater) + references = [] + destroyed = [] + + with patch( + 'Furious.Service.ConnectivityManager.AppSettings.get', return_value=None + ): + for _ in range(30): + manager._testingEnabled = True + reply = _ManagedReply(manager) + references.append(weakref.ref(reply)) + reply.destroyed.connect(lambda *_args: destroyed.append(True)) + requests = [] + + def get(request): + requests.append(request) + + return reply + + with patch.object(manager, 'get', get): + manager.startSingleTest() + manager.startSingleTest() + + self.assertEqual(len(requests), 1) + + deleteQObject(reply) + + self.assertIsNone(manager._activeReply) + self.assertFalse(manager._replyContexts) + self.assertFalse(manager.jobTimeoutTimer.isActive()) + self.assertTrue(manager.jobArrangeTimer.isActive()) + + manager.stopTest() + del reply + + self.assertEqual(len(destroyed), 30) + self.assertTrue(all(reference() is None for reference in references)) + + def testOlderReplyDestructionDoesNotRetireTheCurrentProbe(self): + """A deferred delete from a completed request cannot clear its replacement.""" + manager = ConnectivityManager() + self.addCleanup(manager.deleteLater) + manager._testingEnabled = True + oldReply = _ManagedReply(manager) + currentReply = _ManagedReply(manager) + + with ( + patch.object(manager, 'get', side_effect=[oldReply, currentReply]), + patch( + 'Furious.Service.ConnectivityManager.AppSettings.get', return_value=None + ), + ): + manager.startSingleTest() + oldReply.finished.emit() + manager.startSingleTest() + deleteQObject(oldReply) + + self.assertIs(manager._activeReply, currentReply) + self.assertTrue(manager.jobTimeoutTimer.isActive()) + + manager.stopTest() + + def testStopAfterNativeReplyDestructionDoesNotAccessDeletedWrapper(self): + """Disconnect after reply deletion is harmless even with a retained wrapper.""" + manager = ConnectivityManager() + self.addCleanup(manager.deleteLater) + manager._testingEnabled = True + reply = _ManagedReply(manager) + + with ( + patch.object(manager, 'get', return_value=reply), + patch( + 'Furious.Service.ConnectivityManager.AppSettings.get', return_value=None + ), + ): + manager.startSingleTest() + + deleteQObject(reply) + manager.stopTest() + + self.assertIsNone(manager._activeReply) + self.assertFalse(manager.jobTimeoutTimer.isActive()) + self.assertFalse(manager.jobArrangeTimer.isActive()) + def testRapidStartsReuseOneActiveRequest(self): """Do not accumulate probes or timeout timers during rapid calls.""" manager = ConnectivityManager() - reply = object() + reply = _ManagedReply(manager) manager._testingEnabled = True with ( @@ -599,6 +777,110 @@ class ConnectivityManagerTest(unittest.TestCase): class TrafficStatsManagerTest(unittest.TestCase): """Verify blocked queries and reentrant notifications respect manager lifetime.""" + def testCollectionEnableDoesNotResumeAfterProviderEndsItsOwner(self): + """Monitor discovery may destroy, disconnect, or disable its requester.""" + for action in ('destroy', 'disconnect', 'disable'): + with self.subTest(action=action): + manager = TrafficStatsManager() + manager._connected = True + manager._collectionEnabled = False + monitor = TrafficStatsMonitor( + query=lambda _target: TrafficCounters(1, 2), target=None + ) + + def resolveMonitor(_runtimes): + if action == 'destroy': + deleteQObject(manager) + elif action == 'disconnect': + manager.disconnectedCallback() + else: + manager.setCollectionEnabled(False) + + return monitor + + registry = SimpleNamespace( + trafficStatsMonitorForRuntimes=resolveMonitor + ) + + try: + with patch( + 'Furious.Service.TrafficStatsManager.getPluginRegistry', + return_value=registry, + ): + manager.setCollectionEnabled(True) + + self.assertIsNone(manager._monitor) + self.assertIsNone(manager._executor) + + if isValid(manager): + self.assertFalse(manager._sampleTimer.isActive()) + finally: + if isValid(manager): + manager.cleanup() + deleteQObject(manager) + + def testNativeOwnerDestructionClosesStatisticsExecutor(self): + """Retaining an invalid wrapper must not leave its worker thread running.""" + for blocked in (False, True): + with self.subTest(blocked=blocked): + for _ in range(20): + parent = QtCore.QObject() + manager = TrafficStatsManager(parent) + started = threading.Event() + release = threading.Event() + updates = [] + callbackErrors = [] + + def query(_target): + started.set() + + if blocked: + release.wait(3) + + return TrafficCounters(1, 2) + + manager.sampleChanged.connect(updates.append) + manager._activateMonitor( + TrafficStatsMonitor(query=query, target=None) + ) + executor = manager._executor + + try: + self.assertTrue(started.wait(1)) + threads = tuple(executor._threads) + updates.clear() + + with patch( + 'sys.excepthook', lambda *args: callbackErrors.append(args) + ): + deleteQObject(parent) + + self.assertFalse(isValid(manager)) + self.assertFalse(isValid(manager._sampleTimer)) + self.assertIsNone(manager._executor) + self.assertIsNone(manager._future) + self.assertIsNone(manager._monitor) + self.assertTrue(executor._shutdown) + + release.set() + + for thread in threads: + thread.join(2) + + processQtEvents() + + self.assertFalse(callbackErrors) + self.assertFalse(updates) + self.assertTrue( + all(not thread.is_alive() for thread in threads) + ) + finally: + release.set() + executor.shutdown(wait=True, cancel_futures=True) + + if isValid(parent): + deleteQObject(parent) + def testReconnectPreparationStopsAfterResetOrProviderEndsItsOwner(self): """Do not reactivate statistics after reentrant disconnect or destruction.""" application()