Bound long-lived service retention

Signed-off-by: Loren Eteval <loren.eteval@proton.me>
This commit is contained in:
Loren Eteval
2026-08-23 11:26:08 +08:00
parent 185bc784cb
commit 45f22d3c6f
14 changed files with 373 additions and 49 deletions
+46 -7
View File
@@ -234,6 +234,11 @@ class DesktopApplication(ApplicationRunner, SingletonApplication):
self.mainWindow = None
self.systemTray = None
self._loggingHandlers = tuple()
self._loggingRootLevel = None
self._loggingRaiseExceptions = logging.raiseExceptions
self.connectionController = None
self.routingController = None
self.settingsController = None
@@ -332,16 +337,47 @@ class DesktopApplication(ApplicationRunner, SingletonApplication):
fontFamily=self.customFontName,
)
logging.basicConfig(
format='[%(asctime)s] [%(name)s] [%(levelname)s] %(message)s',
level=logging.INFO,
handlers=(
ApplicationLogHandler(self.logManager),
logging.StreamHandler(),
),
formatter = logging.Formatter(
'[%(asctime)s] [%(name)s] [%(levelname)s] %(message)s'
)
rootLogger = logging.getLogger()
self._loggingRootLevel = rootLogger.level
self._loggingHandlers = (
ApplicationLogHandler(self.logManager),
logging.StreamHandler(),
)
rootLogger.setLevel(logging.INFO)
for handler in self._loggingHandlers:
handler.setFormatter(formatter)
rootLogger.addHandler(handler)
logging.raiseExceptions = False
self._cleanupStack.register('logging', self._cleanupLogging)
def _cleanupLogging(self):
"""Remove and close exactly the handlers owned by this application."""
rootLogger = logging.getLogger()
for handler in self._loggingHandlers:
rootLogger.removeHandler(handler)
handler.close()
self._loggingHandlers = tuple()
if self._loggingRootLevel is not None:
rootLogger.setLevel(self._loggingRootLevel)
self._loggingRootLevel = None
logging.raiseExceptions = self._loggingRaiseExceptions
def addCustomFont(self):
"""Add custom font."""
fontFile = str(DATA_DIR / 'font' / 'CascadiaMono')
@@ -731,6 +767,9 @@ class DesktopApplication(ApplicationRunner, SingletonApplication):
try:
self.addStorage()
except Exception:
# Any non-exit exceptions
# Roll back resources acquired before storage.
Mixins.CleanupOnExit.cleanupAll()
raise
+3 -3
View File
@@ -180,7 +180,7 @@ def classname(ob) -> str:
return ob.__class__.__name__
@functools.lru_cache(None)
@functools.lru_cache(1024)
def isValidIPAddress(address) -> bool:
"""Return whether valid ip address."""
try:
@@ -230,13 +230,13 @@ def runExternalCommand(*args, **kwargs):
return subprocess.run(*args, **kwargs)
@functools.lru_cache(None)
@functools.lru_cache(256)
def absolutePath(path) -> pathlib.Path:
"""Return the absolute path value used by the application."""
return pathlib.Path(path) if os.path.isabs(path) else ROOT_DIR / path
@functools.lru_cache(None)
@functools.lru_cache(128)
def versionToValue(version: str) -> int:
"""Return the version to value value used by the application."""
+10 -6
View File
@@ -41,6 +41,7 @@ class HttpGetManager(AppQNetworkAccessManager):
"""Initialize the HttpGetManager."""
super().__init__(parent)
self.transferTimeout = max(int(kwargs.pop('transferTimeout', 60_000)), 1)
self.actionMessage = actionMessage
self.completionRunsOnce = kwargs.pop('completionRunsOnce', True)
@@ -94,7 +95,7 @@ class HttpGetManager(AppQNetworkAccessManager):
if isinstance(networkReply, QNetworkReply):
self.handleReadyReadByNetworkReply(
networkReply,
**self._replyContexts.get(id(networkReply), {}),
**self._replyContexts.get(networkReply, {}),
)
@QtCore.Slot()
@@ -105,7 +106,7 @@ class HttpGetManager(AppQNetworkAccessManager):
if not isinstance(networkReply, QNetworkReply):
return
kwargs = self._replyContexts.pop(id(networkReply), {})
kwargs = self._replyContexts.pop(networkReply, {})
self.handleFinishedByNetworkReply(networkReply, **kwargs)
@@ -159,13 +160,16 @@ class HttpGetManager(AppQNetworkAccessManager):
def webGET(self, request: Union[QNetworkRequest, str], **kwargs) -> QNetworkReply:
"""Start an HTTP GET request managed by this instance."""
if isinstance(request, QNetworkRequest):
networkReply = self.get(request)
request = QNetworkRequest(request)
else:
networkReply = self.get(QNetworkRequest(QtCore.QUrl(request)))
request = QNetworkRequest(QtCore.QUrl(request))
key = id(networkReply)
if request.transferTimeout() <= 0:
request.setTransferTimeout(self.transferTimeout)
self._replyContexts[key] = dict(kwargs)
networkReply = self.get(request)
self._replyContexts[networkReply] = dict(kwargs)
networkReply.readyRead.connect(self._handleReadyRead)
networkReply.finished.connect(self._handleFinished)
+3 -3
View File
@@ -73,13 +73,13 @@ def setIconAsMask(icon):
return icon
@functools.lru_cache(None)
@functools.lru_cache(256)
def bootstrapIconMask(name):
"""Return the bootstrap icon mask value used by the application."""
return setIconAsMask(bootstrapIconWhite(name))
@functools.lru_cache(None)
@functools.lru_cache(256)
def bootstrapIconWithOpacity(name, opacity, isMask=False):
"""Return the bootstrap icon with opacity value used by the application."""
sourceIcon = bootstrapIconWhite(name)
@@ -223,7 +223,7 @@ class AppQAction(Mixins.QTranslatable, Mixins.ThemeAware, QAction):
return self.textEnglish == compare
@staticmethod
@functools.lru_cache(None)
@functools.lru_cache(256)
def getIconFileName(fileName):
"""Return icon file name."""
try:
+1 -1
View File
@@ -1962,7 +1962,7 @@ class AppQPushButton(Mixins.QTranslatable, Mixins.ThemeAware, QPushButton):
self.setIcon(icon)
@staticmethod
@functools.lru_cache(None)
@functools.lru_cache(256)
def getIconFileName(fileName):
"""Return icon file name."""
try:
+9
View File
@@ -25,6 +25,15 @@
independently; terminal reply paths abort or finish once and schedule deletion once.
- Bound network, DNS, process, host, and worker work where the provider permits it. Executor callbacks use weak or
otherwise bounded ownership and must not retain a manager forever after shutdown; GUI updates cross via signals.
- Log storage has independent count, total-character, and per-entry hard bounds even when automatic clearing is
disabled. High-rate producers must coalesce GUI notifications; hiding a page may defer rendering but must never defer
draining a process pipe or bounded transport queue.
- Metric history has both a time horizon and a defensive sample-count ceiling. Derive graph buckets on demand rather
than retaining a second ever-growing history.
- Every network reply has a finite transfer timeout unless a documented caller supplies a stricter one. Track replies
by exact object identity, remove every context on the terminal path, and schedule the reply for deletion exactly once.
- Per-subscription versions exist only while the persisted subscription, its timer, or an active reply needs them;
repeated create/delete cycles must not grow bookkeeping dictionaries.
- Logging and metrics collect while pages are hidden but avoid hidden-page rendering. Raw time-series samples are
immutable; stable timestamp buckets are derived display data.
- Log retention is owned by `LogManager`, not a text document. Category pruning/reset notifications must keep lazy UI
+120 -11
View File
@@ -57,10 +57,16 @@ class LogManager(QtCore.QObject):
DefaultMaximumEntries = 10_000
DefaultAutoClearMaximumEntries = 5_000
DefaultMaximumCharacters = 8 * 1024 * 1024
DefaultMaximumEntryCharacters = 64 * 1024
TruncationMarker = '\n... [log entry truncated]'
categoryRegistered = QtCore.Signal(object)
entryAdded = QtCore.Signal(object)
entriesCleared = QtCore.Signal(object)
entriesChanged = QtCore.Signal(int)
_entriesChangedRequested = QtCore.Signal()
def __init__(
self,
@@ -69,6 +75,8 @@ class LogManager(QtCore.QObject):
maximumEntries=DefaultMaximumEntries,
autoClearMaximumEntries=DefaultAutoClearMaximumEntries,
autoClearEnabled=True,
maximumCharacters=DefaultMaximumCharacters,
maximumEntryCharacters=DefaultMaximumEntryCharacters,
):
"""Initialize the category registry and thread-safe entry collection."""
super().__init__(parent)
@@ -80,6 +88,23 @@ class LogManager(QtCore.QObject):
):
raise ValueError('maximumEntries must be a positive integer')
if (
isinstance(maximumCharacters, bool)
or not isinstance(maximumCharacters, int)
or maximumCharacters <= 0
):
raise ValueError('maximumCharacters must be a positive integer')
if (
isinstance(maximumEntryCharacters, bool)
or not isinstance(maximumEntryCharacters, int)
or maximumEntryCharacters <= 0
):
raise ValueError('maximumEntryCharacters must be a positive integer')
if maximumEntryCharacters > maximumCharacters:
raise ValueError('maximumEntryCharacters cannot exceed maximumCharacters')
if (
isinstance(autoClearMaximumEntries, bool)
or not isinstance(autoClearMaximumEntries, int)
@@ -90,11 +115,20 @@ class LogManager(QtCore.QObject):
self._lock = threading.RLock()
self._categories: dict[str, LogCategory] = {}
self._maximumEntries = maximumEntries
self._maximumCharacters = maximumCharacters
self._maximumEntryCharacters = maximumEntryCharacters
self._autoClearMaximumEntries = autoClearMaximumEntries
self._autoClearEnabled = bool(autoClearEnabled)
self._entries: deque[LogEntry] = deque(maxlen=maximumEntries)
self._entries: deque[LogEntry] = deque()
self._categoryEntryCounts: dict[str, int] = {}
self._retainedCharacters = 0
self._sequence = 0
self._changeNotificationPending = False
self._entriesChangedRequested.connect(
self._publishEntriesChanged,
QtCore.Qt.ConnectionType.QueuedConnection,
)
self.registerCategory(
LogCategory(
@@ -124,6 +158,22 @@ class LogManager(QtCore.QObject):
"""Return the maximum number of structured entries retained in memory."""
return self._maximumEntries
@property
def maximumCharacters(self) -> int:
"""Return the hard character budget for retained log messages."""
return self._maximumCharacters
@property
def maximumEntryCharacters(self) -> int:
"""Return the maximum number of characters retained from one entry."""
return self._maximumEntryCharacters
@property
def retainedCharacters(self) -> int:
"""Return the current retained message-character count in constant time."""
with self._lock:
return self._retainedCharacters
@property
def autoClearMaximumEntries(self) -> int:
"""Return the automatic-clear threshold for replaceable runtime logs."""
@@ -141,14 +191,73 @@ class LogManager(QtCore.QObject):
def _removeCategoriesLocked(self, categoryIds: set[str]):
"""Remove selected categories while the caller owns ``_lock``."""
self._entries = deque(
(entry for entry in self._entries if entry.categoryId not in categoryIds),
maxlen=self._maximumEntries,
)
retainedEntries = deque()
retainedCharacters = 0
for entry in self._entries:
if entry.categoryId in categoryIds:
continue
retainedEntries.append(entry)
retainedCharacters += len(entry.message)
self._entries = retainedEntries
self._retainedCharacters = retainedCharacters
for categoryId in categoryIds:
self._categoryEntryCounts[categoryId] = 0
def _removeOldestLocked(self):
"""Remove and account for the oldest entry while holding the lock."""
entry = self._entries.popleft()
self._categoryEntryCounts[entry.categoryId] -= 1
self._retainedCharacters -= len(entry.message)
def _enforceRetentionLimitsLocked(self):
"""Evict the oldest entries until every hard retention limit is met."""
while self._entries and (
len(self._entries) > self._maximumEntries
or self._retainedCharacters > self._maximumCharacters
):
self._removeOldestLocked()
def _normalizeMessage(self, message) -> str:
"""Return one newline-trimmed message within the per-entry hard limit."""
text = str(message).rstrip('\r\n')
if len(text) <= self._maximumEntryCharacters:
return text
marker = self.TruncationMarker
if len(marker) >= self._maximumEntryCharacters:
return text[: self._maximumEntryCharacters]
return text[: self._maximumEntryCharacters - len(marker)] + marker
def _requestEntriesChanged(self):
"""Queue at most one cross-thread presentation notification."""
shouldNotify = False
with self._lock:
if not self._changeNotificationPending:
self._changeNotificationPending = True
shouldNotify = True
if shouldNotify:
self._entriesChangedRequested.emit()
@QtCore.Slot()
def _publishEntriesChanged(self):
"""Publish the newest sequence once on the manager's Qt thread."""
with self._lock:
self._changeNotificationPending = False
sequence = self._sequence
self.entriesChanged.emit(sequence)
def setAutoClearEnabled(self, enabled: bool):
"""Apply Core-triggered clearing without involving presentation state."""
clearedCategoryIds = set()
@@ -257,7 +366,7 @@ class LogManager(QtCore.QObject):
self._sequence += 1
entry = LogEntry(
message=str(message).rstrip('\r\n'),
message=self._normalizeMessage(message),
timestamp=timestamp,
categoryId=category.id,
categoryLabel=category.displayName,
@@ -280,19 +389,18 @@ class LogManager(QtCore.QObject):
self._removeCategoriesLocked(clearedCategoryIds)
if len(self._entries) == self._maximumEntries:
evicted = self._entries[0]
self._categoryEntryCounts[evicted.categoryId] -= 1
self._entries.append(entry)
self._categoryEntryCounts[entry.categoryId] += 1
self._retainedCharacters += len(entry.message)
self._enforceRetentionLimitsLocked()
if clearedCategoryIds:
self.entriesCleared.emit(frozenset(clearedCategoryIds))
self.entryAdded.emit(entry)
self._requestEntriesChanged()
return entry
def callback(
@@ -368,6 +476,7 @@ class LogManager(QtCore.QObject):
changed = bool(self._entries)
self._entries.clear()
self._retainedCharacters = 0
for registeredCategoryId in self._categoryEntryCounts:
self._categoryEntryCounts[registeredCategoryId] = 0
+15 -2
View File
@@ -82,6 +82,7 @@ class MetricsHistory(QtCore.QObject):
MaximumHistorySeconds = 24 * 60 * 60
AutoBucketTarget = 120
MaximumSampleCount = 50_000
AutoGranularities = (
1,
2,
@@ -98,7 +99,13 @@ class MetricsHistory(QtCore.QObject):
60 * 60,
)
def __init__(self, parent=None, *, maximumHistorySeconds=None):
def __init__(
self,
parent=None,
*,
maximumHistorySeconds=None,
maximumSampleCount=None,
):
"""Initialize generic metric definitions and bounded sample storage."""
super().__init__(parent)
@@ -107,6 +114,9 @@ class MetricsHistory(QtCore.QObject):
1.0,
)
self._samples = deque()
self._maximumSampleCount = max(
int(maximumSampleCount or self.MaximumSampleCount), 1
)
self._aggregations = {}
self.registerMetric(DOWNLOAD_SPEED_METRIC, MEAN_AGGREGATION)
@@ -242,7 +252,10 @@ class MetricsHistory(QtCore.QObject):
"""Remove samples older than the configured in-memory history."""
oldestAllowed = now - self._maximumHistorySeconds
while self._samples and self._samples[0].sampledAt < oldestAllowed:
while self._samples and (
self._samples[0].sampledAt < oldestAllowed
or len(self._samples) > self._maximumSampleCount
):
self._samples.popleft()
def effectiveGranularity(self, rangeSeconds, granularitySeconds=0) -> float:
+22 -8
View File
@@ -127,6 +127,17 @@ class SubscriptionManager(HttpGetManager):
return version
def _pruneRequestVersion(self, unique: str):
"""Forget version state once no subscription resource needs it."""
if (
unique in Storage.UserSubs()
or unique in self._autoUpdateTimers
or unique in self._replySubscriptions.values()
):
return
self._requestVersions.pop(unique, None)
def _isCurrentRequest(self, kwargs) -> bool:
"""Return whether one completion still targets the current subscription."""
version = kwargs.get('requestVersion')
@@ -230,14 +241,18 @@ class SubscriptionManager(HttpGetManager):
timer.stop()
timer.deleteLater()
self._pruneRequestVersion(unique)
@QtCore.Slot()
def _releaseFinishedReply(self):
"""Forget one exact subscription reply after its completion is dispatched."""
reply = self.sender()
key = id(reply)
unique = self._replySubscriptions.pop(reply, '')
self._activeReplies.pop(key, None)
self._replySubscriptions.pop(key, None)
self._activeReplies.pop(reply, None)
if unique:
self._pruneRequestVersion(unique)
def cancelUpdates(self, unique: str | None = None):
"""Cancel exact active replies and invalidate their eventual completions."""
@@ -252,8 +267,8 @@ class SubscriptionManager(HttpGetManager):
for subscriptionId in subscriptions:
self._nextRequestVersion(subscriptionId)
for key, reply in tuple(self._activeReplies.items()):
if unique is None or self._replySubscriptions.get(key) == unique:
for reply in tuple(self._activeReplies):
if unique is None or self._replySubscriptions.get(reply) == unique:
reply.abort()
def _synchronizeProfiles(self, unique: str, profiles):
@@ -474,10 +489,9 @@ class SubscriptionManager(HttpGetManager):
)
reply = self.webGET(request, logActionMessage=logActionMessage, **kwargs)
key = id(reply)
self._activeReplies[key] = reply
self._replySubscriptions[key] = str(kwargs.get('unique', ''))
self._activeReplies[reply] = reply
self._replySubscriptions[reply] = str(kwargs.get('unique', ''))
reply.finished.connect(self._releaseFinishedReply)
+4 -4
View File
@@ -344,7 +344,7 @@ class LogPage(Mixins.QTranslatable, QMainWindow):
self.autoScrollSwitch.toggled.connect(self._autoScrollChanged)
self.autoClearSwitch.toggled.connect(self._autoClearChanged)
self.manager.categoryRegistered.connect(self._categoryRegistered)
self.manager.entryAdded.connect(self._entryAdded)
self.manager.entriesChanged.connect(self._entriesChanged)
self.manager.entriesCleared.connect(self._entriesCleared)
scrollbar = self.textBrowser.verticalScrollBar()
@@ -782,10 +782,10 @@ class LogPage(Mixins.QTranslatable, QMainWindow):
self.filterComboBox.findData(category.id)
)
@QtCore.Slot(object)
def _entryAdded(self, entry):
@QtCore.Slot(int)
def _entriesChanged(self, sequence: int):
"""Coalesce newly collected entries into the visible presentation."""
if entry.sequence <= self._renderedSequence:
if sequence <= self._renderedSequence:
return
self._requestRefresh()
+18
View File
@@ -283,6 +283,24 @@ class FrozenlibUtilityTest(unittest.TestCase):
('2001:db8::1', '8443'),
)
def testExternallyKeyedUtilityCachesAreBounded(self):
"""Prevent address, path, and version inputs from growing global caches."""
cachedFunctions = (
UtilityModule.isValidIPAddress,
UtilityModule.absolutePath,
UtilityModule.versionToValue,
)
for function in cachedFunctions:
function.cache_clear()
for index in range(function.cache_info().maxsize + 25):
function(str(index))
info = function.cache_info()
self.assertIsNotNone(info.maxsize)
self.assertLessEqual(info.currsize, info.maxsize)
def testTcpingUsesResolverIPv6AndFallsBackToIPv4(self):
"""Use getaddrinfo candidates rather than IPv4-only gethostbyname."""
candidates = [
+42
View File
@@ -450,6 +450,34 @@ class LogManagerTest(unittest.TestCase):
self.assertEqual(manager.entryCount(TUN2SOCKS_LOG_CATEGORY), 0)
self.assertEqual(manager.entryCount(APPLICATION_LOG_CATEGORY), 1)
def testCharacterBudgetsRemainHardWhenAutoClearIsDisabled(self):
"""Bound both total text and one hostile entry independently of count."""
manager = LogManager(
maximumEntries=100,
maximumCharacters=80,
maximumEntryCharacters=32,
autoClearEnabled=False,
)
manager.append('a' * 100, CORE_LOG_CATEGORY)
self.assertEqual(len(manager.entries()[0].message), 32)
self.assertIn('truncated', manager.entries()[0].message)
for index in range(20):
manager.append(f'{index:02d}-' + ('x' * 17), CORE_LOG_CATEGORY)
self.assertLessEqual(manager.retainedCharacters, 80)
self.assertLessEqual(manager.entryCount(), 4)
self.assertEqual(
manager.retainedCharacters,
sum(len(entry.message) for entry in manager.entries()),
)
self.assertEqual(
manager.entryCount(CORE_LOG_CATEGORY),
manager.entryCount(),
)
class MetricsHistoryTest(unittest.TestCase):
"""Verify bounded history and metric-specific aggregation semantics."""
@@ -522,6 +550,20 @@ class MetricsHistoryTest(unittest.TestCase):
self.assertEqual(manager.rawSamples(), tuple())
self.assertEqual(changed, [])
def testDefensiveSampleCeilingBoundsBurstCadence(self):
"""Retain only the newest samples even when time does not advance."""
manager = MetricsHistory(
maximumHistorySeconds=24 * 60 * 60,
maximumSampleCount=25,
)
for index in range(1000):
manager.recordSample({DOWNLOAD_SPEED_METRIC: index}, sampledAt=1)
self.assertEqual(manager.sampleCount(), 25)
self.assertEqual(manager.rawSamples()[0].values[DOWNLOAD_SPEED_METRIC], 975)
self.assertEqual(manager.rawSamples()[-1].values[DOWNLOAD_SPEED_METRIC], 999)
class TranslationExtractorTest(unittest.TestCase):
"""Verify static translation extraction and constant interpolation."""
+56 -2
View File
@@ -26,11 +26,13 @@ from Furious.Plugins import (
TrafficStatsMonitor,
)
from Furious.Service.ConnectivityManager import ConnectivityManager
from Furious.Qt.HttpGetManager import HttpGetManager
from Furious.Service.PluginUIManager import PluginNavigationManager
from Furious.Service.TrafficStatsManager import TrafficStatsManager
from Furious.Service.UpdateManager import UpdateManager
from PySide6 import QtCore
from PySide6.QtNetwork import QNetworkReply
from PySide6.QtWidgets import QWidget
from shiboken6 import isValid
@@ -117,11 +119,63 @@ class UpdateManagerTest(unittest.TestCase):
parent.deleteLater()
class _NavigationProvider:
"""Return one valid page and one invalid parented QObject."""
class _ManagedReply(QNetworkReply):
"""Provide a hermetic reply object with real Qt lifecycle signals."""
def __init__(self, parent=None):
super().__init__(parent)
self.open(QtCore.QIODevice.OpenModeFlag.ReadOnly)
def abort(self):
self.setFinished(True)
def readData(self, maximumLength):
return bytes()
class _CapturingHttpGetManager(HttpGetManager):
"""Capture normalized requests without performing network I/O."""
def __init__(self):
super().__init__()
self.reply = _ManagedReply(self)
self.request = None
def get(self, request):
self.request = request
return self.reply
class HttpGetManagerLifetimeTest(unittest.TestCase):
"""Verify every request receives a timeout and releases exact context."""
@classmethod
def setUpClass(cls):
application()
def testRequestHasFiniteTimeoutAndTerminalPathDropsContext(self):
manager = _CapturingHttpGetManager()
reply = manager.webGET('https://invalid.test/resource', marker='fixture')
self.assertIs(reply, manager.reply)
self.assertEqual(manager.request.transferTimeout(), 60_000)
self.assertEqual(manager._replyContexts[reply], {'marker': 'fixture'})
with patch.object(manager, 'handleFinishedByNetworkReply') as finished:
reply.finished.emit()
self.assertEqual(manager._replyContexts, {})
finished.assert_called_once_with(reply, marker='fixture')
manager.deleteLater()
capabilityId = 'fixture.navigation'
class _NavigationProvider:
"""Return one valid page and one invalid parented QObject."""
def __init__(self):
"""Initialize construction counters used by idempotence assertions."""
self.validCalls = 0
+24 -2
View File
@@ -214,10 +214,10 @@ class SubscriptionManagerTest(TestCase):
groupBReply = _AbortableReply()
manager._requestVersions.update({'group-a': 1, 'group-b': 4})
manager._activeReplies.update(
{id(groupAReply): groupAReply, id(groupBReply): groupBReply}
{groupAReply: groupAReply, groupBReply: groupBReply}
)
manager._replySubscriptions.update(
{id(groupAReply): 'group-a', id(groupBReply): 'group-b'}
{groupAReply: 'group-a', groupBReply: 'group-b'}
)
manager.cancelUpdates('group-a')
@@ -264,3 +264,25 @@ class SubscriptionManagerTest(TestCase):
self.assertEqual(manager._autoUpdateTimers, {})
self.assertFalse(timer.isActive())
manager.deleteLater()
def testDeletedSubscriptionVersionIsPrunedAfterItsReplyFinishes(self):
"""Release stale-request bookkeeping after the last exact owner ends."""
manager = self._manager()
reply = _AbortableReply()
manager._requestVersions['deleted-group'] = 4
manager._activeReplies[reply] = reply
manager._replySubscriptions[reply] = 'deleted-group'
with mock.patch(
'Furious.Service.SubscriptionManager.Storage.UserSubs',
return_value={},
):
manager._pruneRequestVersion('deleted-group')
self.assertIn('deleted-group', manager._requestVersions)
manager._activeReplies.pop(reply)
manager._replySubscriptions.pop(reply)
manager._pruneRequestVersion('deleted-group')
self.assertNotIn('deleted-group', manager._requestVersions)
manager.deleteLater()