mirror of
https://github.com/LorenEteval/Furious.git
synced 2026-10-08 22:59:48 +03:00
Fix Qt service lifetime cleanup
Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
+169
-1
@@ -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(
|
||||
|
||||
@@ -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)]
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user