Fix subscription reply cleanup

Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
Loren Eteval
2026-10-06 12:11:23 +08:00
parent 7e168e5d4c
commit fe40bec14d
2 changed files with 90 additions and 13 deletions
+24 -13
View File
@@ -42,6 +42,8 @@ from Furious.Service.SubscriptionSync import SubscriptionSynchronizer
from PySide6 import QtCore
from PySide6.QtNetwork import QNetworkRequest
from shiboken6 import isValid
from dataclasses import dataclass
import re
@@ -420,7 +422,10 @@ class SubscriptionManager(HttpGetManager):
for reply in tuple(self._activeReplies):
if unique is None or self._replySubscriptions.get(reply) == unique:
reply.abort()
# Aborting another reply can synchronously destroy this manager
# and every remaining reply in the snapshot.
if isValid(reply):
reply.abort()
for job in tuple(self._preparationJobs.values()):
if unique is None or job.context.get('unique') == unique:
@@ -480,18 +485,22 @@ class SubscriptionManager(HttpGetManager):
timer.stop()
self.cancelUpdates()
self._preparationPool.clear()
# The timeout is diagnostic, not permission to destroy running workers.
# Keep the relay/pool alive until every worker has actually finished.
if not self._preparationPool.waitForDone(self.ShutdownWarningMilliseconds):
logger.warning(
'subscription preparation has not stopped after %s ms; '
'waiting for running workers to finish',
self.ShutdownWarningMilliseconds,
)
# Reply abort hooks can destroy the manager and its child pool. Native
# pool destruction already waits for its workers; do not use that wrapper.
if isValid(self._preparationPool):
self._preparationPool.clear()
self._preparationPool.waitForDone()
# The timeout is diagnostic, not permission to destroy running workers.
# Keep the relay/pool alive until every worker has actually finished.
if not self._preparationPool.waitForDone(self.ShutdownWarningMilliseconds):
logger.warning(
'subscription preparation has not stopped after %s ms; '
'waiting for running workers to finish',
self.ShutdownWarningMilliseconds,
)
self._preparationPool.waitForDone()
self._preparationJobs.clear()
self._preparationPayloads.clear()
@@ -1201,8 +1210,10 @@ class SubscriptionManager(HttpGetManager):
reply = self.webGET(request, logActionMessage=logActionMessage, **kwargs)
self._activeReplies[reply] = reply
self._replySubscriptions[reply] = str(kwargs.get('unique', ''))
self._trackReplyContext(reply, self._activeReplies, reply)
self._trackReplyContext(
reply, self._replySubscriptions, str(kwargs.get('unique', ''))
)
connectWeakly(
reply.finished,
+66
View File
@@ -29,6 +29,8 @@ from Furious.Service.ConnectivityManager import ConnectivityManager
from Furious.Service.EndpointInfoService import ProxyEndpointHttpClient
from Furious.Qt.HttpGetManager import HttpGetManager
from Furious.Service.PluginUIManager import PluginNavigationManager
from Furious.Service.SubscriptionManager import SubscriptionManager
from Furious.Repository import Storage
from Furious.Service.TrafficStatsManager import TrafficStatsManager
from Furious.Service.UpdateManager import UpdateManager
@@ -225,6 +227,70 @@ class HttpGetManagerLifetimeTest(unittest.TestCase):
self.assertTrue(all(not isValid(reply) for reply in replies))
self.assertEqual(manager._pendingRequests, {})
def testSubscriptionRepliesReleaseTrackingOnNativeDestruction(self):
"""Early deletion releases all reply keys before later cancellation."""
with patch.object(Storage, 'UserSubs', return_value={}):
manager = SubscriptionManager()
references = []
destroyed = []
try:
for index in range(30):
reply = _ManagedReply(manager)
references.append(weakref.ref(reply))
reply.destroyed.connect(lambda *_args: destroyed.append(True))
with patch.object(manager, 'get', return_value=reply):
manager.updateSubsByWebGET(
webURL='https://invalid.test', unique=str(index)
)
self.assertIn(reply, manager._activeReplies)
reply.deleteLater()
processQtEvents()
self.assertFalse(isValid(reply))
del reply
self.assertFalse(manager._replyContexts)
self.assertFalse(manager._activeReplies)
self.assertFalse(manager._replySubscriptions)
manager.cancelUpdates()
self.assertEqual(len(destroyed), 30)
collectAtBoundary()
self.assertTrue(all(reference() is None for reference in references))
finally:
manager.shutdown()
manager.deleteLater()
processQtEvents()
def testSubscriptionCancellationToleratesOwnerDestructionDuringAbort(self):
"""Reentrant destruction invalidates every remaining snapshot reply."""
for method in ('cancelUpdates', 'shutdown'):
with self.subTest(method=method):
self._cancelSubscriptionWithReentrantDestruction(method)
def _cancelSubscriptionWithReentrantDestruction(self, method):
"""Keep both public teardown entry points on the real native boundary."""
with patch.object(Storage, 'UserSubs', return_value={}):
manager = SubscriptionManager()
replies = [_ManagedReply(manager), _ManagedReply(manager)]
for index, reply in enumerate(replies):
with patch.object(manager, 'get', return_value=reply):
manager.updateSubsByWebGET(
webURL='https://invalid.test', unique=str(index)
)
replies[0].abort = lambda: deleteQObject(manager)
getattr(manager, method)()
self.assertFalse(isValid(manager))
self.assertTrue(all(not isValid(reply) for reply in replies))
self.assertFalse(manager._replyContexts)
self.assertFalse(manager._activeReplies)
self.assertFalse(manager._replySubscriptions)
def testRequestHasFiniteTimeoutAndTerminalPathDropsContext(self):
manager = _CapturingHttpGetManager()