From fbd28c48d4a602750f67e103c8335e797c2766c3 Mon Sep 17 00:00:00 2001 From: Loren Eteval Date: Wed, 26 Aug 2026 23:30:26 +0800 Subject: [PATCH] Implement generation-based log storage --- Furious/Service/LogManager.py | 503 +++++++-- Furious/Window/LogPage.py | 6 + tests/README.md | 4 + tests/test_log_manager_generation.py | 1430 ++++++++++++++++++++++++++ tests/test_models_and_services.py | 737 ++++++++++++- 5 files changed, 2596 insertions(+), 84 deletions(-) create mode 100644 tests/test_log_manager_generation.py diff --git a/Furious/Service/LogManager.py b/Furious/Service/LogManager.py index 166b9090..00b61d28 100644 --- a/Furious/Service/LogManager.py +++ b/Furious/Service/LogManager.py @@ -23,8 +23,9 @@ from Furious.Models.Logging import LogCategory, LogEntry from PySide6 import QtCore -from collections import deque +from collections import OrderedDict, deque from datetime import datetime +from heapq import merge from typing import Optional import logging @@ -52,13 +53,142 @@ def formatLogEntry(entry: LogEntry) -> str: return entry.message +class _CategoryEntries: + """Own one category's ordered entries and aggregate character count.""" + + __slots__ = ('entries', 'characterCount') + + def __init__(self): + """Initialize an empty category index.""" + self.entries: OrderedDict[int, LogEntry] = OrderedDict() + self.characterCount = 0 + + @property + def entryCount(self) -> int: + """Return the number of physically owned entries.""" + return len(self.entries) + + def append(self, entry: LogEntry, characterCount: int): + """Append one entry and update aggregate accounting.""" + self.entries[entry.sequence] = entry + self.characterCount += characterCount + + def remove(self, sequence: int) -> LogEntry: + """Remove one known sequence and update aggregate accounting.""" + entry = self.entries.pop(sequence) + self.characterCount -= len(entry.message) + + return entry + + def discardOldestPhysical(self) -> int: + """Drop one retired entry and return its released character count.""" + _sequence, entry = self.entries.popitem(last=False) + characterCount = len(entry.message) + self.characterCount -= characterCount + + return characterCount + + +class _EntryGeneration: + """Own one active or retired chronological generation of log entries.""" + + __slots__ = ( + 'identifier', + 'scope', + 'entries', + 'entriesByCategory', + 'characterCount', + ) + + def __init__(self, identifier: int, scope: str): + """Initialize an empty generation for one fixed retention scope.""" + self.identifier = identifier + self.scope = scope + self.entries: OrderedDict[int, LogEntry] = OrderedDict() + self.entriesByCategory: dict[str, _CategoryEntries] = {} + self.characterCount = 0 + + @property + def entryCount(self) -> int: + """Return the number of physically owned entries.""" + return len(self.entries) + + def categoryEntries(self, categoryId: str) -> Optional[_CategoryEntries]: + """Return the active index for one category when it has entries.""" + return self.entriesByCategory.get(categoryId) + + def append(self, entry: LogEntry, characterCount: int): + """Append the same immutable entry to global and category indexes.""" + categoryEntries = self.entriesByCategory.get(entry.categoryId) + + if categoryEntries is None: + categoryEntries = _CategoryEntries() + self.entriesByCategory[entry.categoryId] = categoryEntries + + self.entries[entry.sequence] = entry + categoryEntries.append(entry, characterCount) + self.characterCount += characterCount + + def removeOldest(self) -> LogEntry: + """Remove the oldest entry from both live indexes in constant time.""" + sequence, entry = self.entries.popitem(last=False) + categoryEntries = self.entriesByCategory[entry.categoryId] + indexedEntry = categoryEntries.remove(sequence) + + if indexedEntry is not entry: + raise RuntimeError('log generation indexes diverged') + + if not categoryEntries.entryCount: + del self.entriesByCategory[entry.categoryId] + + self.characterCount -= len(entry.message) + + return entry + + def detachCategory(self, categoryId: str) -> Optional[_CategoryEntries]: + """Detach one category while preserving its entries for deferred release. + + Selective deletion must unlink each sequence from this generation's flat + chronological index. The detached category index retains the LogEntry + objects so their final decrefs can still be performed in bounded batches. + """ + categoryEntries = self.entriesByCategory.pop(categoryId, None) + + if categoryEntries is None: + return None + + for sequence in categoryEntries.entries: + del self.entries[sequence] + + self.characterCount -= categoryEntries.characterCount + + return categoryEntries + + def discardOldestPhysical(self) -> int: + """Drop one entry from a retired generation's duplicate indexes.""" + entry = self.removeOldest() + + return len(entry.message) + + class LogManager(QtCore.QObject): - """Own the application-wide categorized log stream.""" + """Own globally ordered logs through small active generation streams. + + Application and runtime entries deliberately share the same hard count and + character budgets. Generation rollover makes runtime-wide logical clearing + independent of retained history size, while retired object destruction is + spread across bounded cleanup batches. + """ DefaultMaximumEntries = 10_000 DefaultAutoClearMaximumEntries = 5_000 DefaultMaximumCharacters = 8 * 1024 * 1024 DefaultMaximumEntryCharacters = 64 * 1024 + # At least one stale entry is reclaimed before every new entry is retained. + # A larger fixed budget lets normal traffic drain old generations quickly; + # because a clear cannot retire more live entries than maximumEntries, this + # also bounds the live-plus-retired entry backlog under repeated fast clears. + RetiredCleanupBudget = 64 TruncationMarker = '\n... [log entry truncated]' categoryRegistered = QtCore.Signal(object) @@ -67,6 +197,7 @@ class LogManager(QtCore.QObject): entriesChanged = QtCore.Signal(int) _entriesChangedRequested = QtCore.Signal() + _retiredCleanupRequested = QtCore.Signal() def __init__( self, @@ -119,16 +250,41 @@ class LogManager(QtCore.QObject): self._maximumEntryCharacters = maximumEntryCharacters self._autoClearMaximumEntries = autoClearMaximumEntries self._autoClearEnabled = bool(autoClearEnabled) - self._entries: deque[LogEntry] = deque() - self._categoryEntryCounts: dict[str, int] = {} - self._retainedCharacters = 0 self._sequence = 0 + self._retainedEntryCount = 0 + self._retainedCharacters = 0 + self._generationSequence = 0 + # Three fixed active streams preserve exact global ordering with an O(n) + # three-way merge and make oldest-live selection constant time. Runtime + # categories are isolated so disconnect/runtime clear is one reference + # swap; other non-Application categories use a second generation because + # Core auto-clear historically clears them too, while runtimeOnly does not. + self._applicationGeneration = self._newGenerationLocked('application') + self._runtimeGeneration = self._newGenerationLocked('runtime') + self._otherGeneration = self._newGenerationLocked('other') + self._runtimeCategoryIds = frozenset() + self._nonApplicationCategoryIds = frozenset() + # Retired batches remain strongly owned until bounded cleanup removes + # their entries. Merely dropping a large OrderedDict here would make the + # clear caller synchronously execute thousands of CPython decrefs. + self._retiredBatches = deque() + self._retiredEntryCount = 0 + self._retiredCharacters = 0 + self._retiredCleanupPending = False self._changeNotificationPending = False + # LogManager is one application-lifetime QObject parented to the desktop + # application in production. These are intentional queued self-connections: + # sender and receiver share one destruction boundary, no transient UI is + # retained, and pending delivery is discarded when that QObject is destroyed. self._entriesChangedRequested.connect( self._publishEntriesChanged, QtCore.Qt.ConnectionType.QueuedConnection, ) + self._retiredCleanupRequested.connect( + self._cleanupRetired, + QtCore.Qt.ConnectionType.QueuedConnection, + ) self.registerCategory( LogCategory( @@ -174,6 +330,18 @@ class LogManager(QtCore.QObject): with self._lock: return self._retainedCharacters + @property + def retiredEntryCount(self) -> int: + """Return logically dead entries awaiting incremental physical release.""" + with self._lock: + return self._retiredEntryCount + + @property + def retiredCharacters(self) -> int: + """Return message characters awaiting incremental physical release.""" + with self._lock: + return self._retiredCharacters + @property def autoClearMaximumEntries(self) -> int: """Return the automatic-clear threshold for replaceable runtime logs.""" @@ -185,40 +353,147 @@ class LogManager(QtCore.QObject): with self._lock: return self._autoClearEnabled - def _nonApplicationCategoryIdsLocked(self) -> set[str]: - """Return every registered category except the persistent Application log.""" - return set(self._categories).difference({APPLICATION_LOG_CATEGORY}) + def _newGenerationLocked(self, scope: str) -> _EntryGeneration: + """Create one uniquely identified empty active generation.""" + self._generationSequence += 1 - def _removeCategoriesLocked(self, categoryIds: set[str]): - """Remove selected categories while the caller owns ``_lock``.""" - retainedEntries = deque() - retainedCharacters = 0 + return _EntryGeneration(self._generationSequence, scope) - for entry in self._entries: - if entry.categoryId in categoryIds: - continue + def _activeGenerationsLocked(self) -> tuple[_EntryGeneration, ...]: + """Return the fixed active streams in merge-independent order.""" + return ( + self._applicationGeneration, + self._runtimeGeneration, + self._otherGeneration, + ) - retainedEntries.append(entry) - retainedCharacters += len(entry.message) + def _liveEntryCountLocked(self) -> int: + """Return the incrementally maintained total live entry count.""" + return self._retainedEntryCount - self._entries = retainedEntries - self._retainedCharacters = retainedCharacters + def _liveCharactersLocked(self) -> int: + """Return the incrementally maintained total live character count.""" + return self._retainedCharacters - for categoryId in categoryIds: - self._categoryEntryCounts[categoryId] = 0 + @staticmethod + def _generationAttributeForCategory(category: LogCategory) -> str: + """Return the active generation attribute that owns one category.""" + if category.id == APPLICATION_LOG_CATEGORY: + return '_applicationGeneration' + + return '_runtimeGeneration' if category.runtime else '_otherGeneration' + + def _generationForCategoryLocked(self, category: LogCategory) -> _EntryGeneration: + """Return the active generation that owns one registered category.""" + if category.id == APPLICATION_LOG_CATEGORY: + return self._applicationGeneration + + return self._runtimeGeneration if category.runtime else self._otherGeneration + + def _retireBatchLocked(self, batch): + """Keep one logically dead non-empty batch for bounded reclamation.""" + if not batch.entryCount: + return + + self._retiredBatches.append(batch) + self._retiredEntryCount += batch.entryCount + self._retiredCharacters += batch.characterCount + + def _rollGenerationLocked(self, attributeName: str) -> int: + """Swap one active generation and retire its old entries in O(1).""" + generation = getattr(self, attributeName) + entryCount = generation.entryCount + characterCount = generation.characterCount + setattr( + self, + attributeName, + self._newGenerationLocked(generation.scope), + ) + self._retainedEntryCount -= entryCount + self._retainedCharacters -= characterCount + self._retireBatchLocked(generation) + + return entryCount + + def _clearNonApplicationLocked(self) -> bool: + """Logically clear every non-Application stream with two swaps.""" + removedEntries = self._rollGenerationLocked('_runtimeGeneration') + removedEntries += self._rollGenerationLocked('_otherGeneration') + + return bool(removedEntries) + + def _cleanupRetiredLocked(self, budget: Optional[int] = None) -> int: + """Physically release at most one fixed batch of logically dead entries.""" + remaining = ( + max(1, int(self.RetiredCleanupBudget)) + if budget is None + else max(0, int(budget)) + ) + cleaned = 0 + + while remaining and self._retiredBatches: + batch = self._retiredBatches[0] + characterCount = batch.discardOldestPhysical() + self._retiredEntryCount -= 1 + self._retiredCharacters -= characterCount + cleaned += 1 + remaining -= 1 + + if not batch.entryCount: + if batch.characterCount: + raise RuntimeError('retired log accounting diverged') + + self._retiredBatches.popleft() + + return cleaned + + def _requestRetiredCleanup(self): + """Queue one bounded reclamation turn when retired entries remain.""" + shouldNotify = False + + with self._lock: + if self._retiredBatches and not self._retiredCleanupPending: + self._retiredCleanupPending = True + shouldNotify = True + + if shouldNotify: + self._retiredCleanupRequested.emit() + + @QtCore.Slot() + def _cleanupRetired(self): + """Release one bounded retired batch on the manager's Qt thread.""" + with self._lock: + self._retiredCleanupPending = False + self._cleanupRetiredLocked() + shouldContinue = bool(self._retiredBatches) + + if shouldContinue: + self._requestRetiredCleanup() def _removeOldestLocked(self): - """Remove and account for the oldest entry while holding the lock.""" - entry = self._entries.popleft() + """Remove the oldest live entry across three streams in constant time.""" + generations = tuple( + generation + for generation in self._activeGenerationsLocked() + if generation.entryCount + ) - self._categoryEntryCounts[entry.categoryId] -= 1 + if not generations: + raise RuntimeError('cannot evict from an empty log manager') + + oldestGeneration = min( + generations, + key=lambda generation: next(iter(generation.entries)), + ) + entry = oldestGeneration.removeOldest() + self._retainedEntryCount -= 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 + while self._liveEntryCountLocked() and ( + self._liveEntryCountLocked() > self._maximumEntries + or self._liveCharactersLocked() > self._maximumCharacters ): self._removeOldestLocked() @@ -236,18 +511,14 @@ class LogManager(QtCore.QObject): return text[: self._maximumEntryCharacters - len(marker)] + marker - def _requestEntriesChanged(self): - """Queue at most one cross-thread presentation notification.""" - shouldNotify = False + def _markEntriesChangedLocked(self) -> bool: + """Mark one coalesced presentation notification while already locked.""" + if self._changeNotificationPending: + return False - with self._lock: - if not self._changeNotificationPending: - self._changeNotificationPending = True + self._changeNotificationPending = True - shouldNotify = True - - if shouldNotify: - self._entriesChangedRequested.emit() + return True @QtCore.Slot() def _publishEntriesChanged(self): @@ -260,33 +531,41 @@ class LogManager(QtCore.QObject): def setAutoClearEnabled(self, enabled: bool): """Apply Core-triggered clearing without involving presentation state.""" - clearedCategoryIds = set() + clearedCategoryIds = frozenset() with self._lock: self._autoClearEnabled = bool(enabled) if ( self._autoClearEnabled - and self._categoryEntryCounts.get(CORE_LOG_CATEGORY, 0) + and self._categoryEntryCountLocked(CORE_LOG_CATEGORY) >= self._autoClearMaximumEntries ): - clearedCategoryIds = self._nonApplicationCategoryIdsLocked() - - self._removeCategoriesLocked(clearedCategoryIds) + clearedCategoryIds = self._nonApplicationCategoryIds + self._clearNonApplicationLocked() if clearedCategoryIds: - self.entriesCleared.emit(frozenset(clearedCategoryIds)) + self._requestRetiredCleanup() + self.entriesCleared.emit(clearedCategoryIds) + + def _categoryEntryCountLocked(self, categoryId: str) -> int: + """Return one registered category's live count in constant time.""" + category = self._categories[categoryId] + generation = self._generationForCategoryLocked(category) + categoryEntries = generation.categoryEntries(categoryId) + + return categoryEntries.entryCount if categoryEntries is not None else 0 def entryCount(self, categoryId: Optional[str] = None) -> int: """Return the retained total or one category count in constant time.""" with self._lock: if categoryId in (None, ALL_LOGS_FILTER): - return len(self._entries) + return self._liveEntryCountLocked() if categoryId not in self._categories: raise KeyError(f'unknown log category {categoryId!r}') - return self._categoryEntryCounts.get(categoryId, 0) + return self._categoryEntryCountLocked(categoryId) def registerCategory(self, category: LogCategory) -> LogCategory: """Register a filterable category and publish it exactly once.""" @@ -306,7 +585,14 @@ class LogManager(QtCore.QObject): return existing self._categories[category.id] = category - self._categoryEntryCounts[category.id] = 0 + + if category.id != APPLICATION_LOG_CATEGORY: + self._nonApplicationCategoryIds = self._nonApplicationCategoryIds.union( + {category.id} + ) + + if category.runtime: + self._runtimeCategoryIds = self._runtimeCategoryIds.union({category.id}) self.categoryRegistered.emit(category) @@ -349,8 +635,10 @@ class LogManager(QtCore.QObject): source: str = '', severity: str = '', ) -> LogEntry: - """Append a structured entry and notify interested presenters.""" - clearedCategoryIds = set() + """Append with fixed cleanup work and notify interested presenters.""" + clearedCategoryIds = frozenset() + shouldNotifyEntriesChanged = False + retiredCleanupNeeded = False with self._lock: category = self._categories.get(categoryId) @@ -358,6 +646,12 @@ class LogManager(QtCore.QObject): if category is None: raise KeyError(f'unknown log category {categoryId!r}') + # Reclamation is deliberately entry-budgeted rather than generation- + # budgeted. Even if a previous clear retired thousands of objects, + # this append performs at most RetiredCleanupBudget physical releases. + if self._retiredBatches: + self._cleanupRetiredLocked() + if timestamp is None: timestamp = datetime.now() elif not isinstance(timestamp, datetime): @@ -379,27 +673,35 @@ class LogManager(QtCore.QObject): if ( category.id == CORE_LOG_CATEGORY and self._autoClearEnabled - and self._categoryEntryCounts[CORE_LOG_CATEGORY] + and self._categoryEntryCountLocked(CORE_LOG_CATEGORY) >= self._autoClearMaximumEntries ): # When the next Core line would exceed its threshold, retain # only Application diagnostics. All other categories belong - # to the replaceable runtime stream for this policy. - clearedCategoryIds = self._nonApplicationCategoryIdsLocked() + # to two replaceable streams that roll over without touching + # their entries or performing synchronous object destruction. + clearedCategoryIds = self._nonApplicationCategoryIds + self._clearNonApplicationLocked() - self._removeCategoriesLocked(clearedCategoryIds) - - self._entries.append(entry) - self._categoryEntryCounts[entry.categoryId] += 1 - self._retainedCharacters += len(entry.message) + characterCount = len(entry.message) + generation = self._generationForCategoryLocked(category) + generation.append(entry, characterCount) + self._retainedEntryCount += 1 + self._retainedCharacters += characterCount self._enforceRetentionLimitsLocked() + shouldNotifyEntriesChanged = self._markEntriesChangedLocked() + retiredCleanupNeeded = bool(self._retiredBatches) if clearedCategoryIds: - self.entriesCleared.emit(frozenset(clearedCategoryIds)) + self.entriesCleared.emit(clearedCategoryIds) self.entryAdded.emit(entry) - self._requestEntriesChanged() + if shouldNotifyEntriesChanged: + self._entriesChangedRequested.emit() + + if retiredCleanupNeeded: + self._requestRetiredCleanup() return entry @@ -436,14 +738,34 @@ class LogManager(QtCore.QObject): self, categoryId: Optional[str] = None, ) -> tuple[int, tuple[LogEntry, ...]]: - """Return the current sequence and its immutable filtered entries.""" + """Return an O(n) global or O(k) category-specific immutable snapshot.""" with self._lock: if categoryId in (None, ALL_LOGS_FILTER): - entries = tuple(self._entries) - else: + # Exactly three sorted active streams make heap merge O(n) with + # a constant fan-in. Retired generations are never visited, so + # logically cleared entries cannot tax or leak into snapshots. entries = tuple( - entry for entry in self._entries if entry.categoryId == categoryId + merge( + *( + generation.entries.values() + for generation in self._activeGenerationsLocked() + ), + key=lambda entry: entry.sequence, + ) ) + else: + category = self._categories.get(categoryId) + + if category is None: + entries = tuple() + else: + generation = self._generationForCategoryLocked(category) + categoryEntries = generation.categoryEntries(categoryId) + entries = ( + tuple(categoryEntries.entries.values()) + if categoryEntries is not None + else tuple() + ) return self._sequence, entries @@ -453,43 +775,58 @@ class LogManager(QtCore.QObject): *, runtimeOnly: bool = False, ): - """Clear all, one category, or every transient runtime category.""" + """Logically clear a generation or selectively unlink one category. + + Whole-history, Application-only, runtime-only, and sole-category clears + are O(1) generation swaps. A category sharing a generation with other + live categories requires O(k) keyed unlinks to keep the active stream + tombstone-free for exact oldest-live eviction and O(n) global snapshots. + Its LogEntry destruction is still deferred in bounded cleanup batches. + """ if categoryId is not None and runtimeOnly: raise ValueError('categoryId and runtimeOnly cannot be combined') with self._lock: if runtimeOnly: - clearedCategoryIds = { - category.id - for category in self._categories.values() - if category.runtime - } + clearedCategoryIds = self._runtimeCategoryIds + changed = bool(self._rollGenerationLocked('_runtimeGeneration')) elif categoryId is None or categoryId == ALL_LOGS_FILTER: clearedCategoryIds = None + changed = bool(self._liveEntryCountLocked()) + + self._rollGenerationLocked('_applicationGeneration') + self._rollGenerationLocked('_runtimeGeneration') + self._rollGenerationLocked('_otherGeneration') else: - if categoryId not in self._categories: + category = self._categories.get(categoryId) + + if category is None: raise KeyError(f'unknown log category {categoryId!r}') - clearedCategoryIds = {categoryId} + clearedCategoryIds = frozenset({categoryId}) + attributeName = self._generationAttributeForCategory(category) + generation = getattr(self, attributeName) + categoryEntries = generation.categoryEntries(categoryId) + changed = categoryEntries is not None - if clearedCategoryIds is None: - changed = bool(self._entries) + if changed and categoryEntries.entryCount == generation.entryCount: + self._rollGenerationLocked(attributeName) + elif changed: + retiredCategory = generation.detachCategory(categoryId) - self._entries.clear() - self._retainedCharacters = 0 + if retiredCategory is None: + raise RuntimeError( + 'log category index disappeared during clear' + ) - for registeredCategoryId in self._categoryEntryCounts: - self._categoryEntryCounts[registeredCategoryId] = 0 - else: - oldLength = len(self._entries) - - self._removeCategoriesLocked(clearedCategoryIds) - - changed = len(self._entries) != oldLength + self._retainedEntryCount -= retiredCategory.entryCount + self._retainedCharacters -= retiredCategory.characterCount + self._retireBatchLocked(retiredCategory) if changed: + self._requestRetiredCleanup() self.entriesCleared.emit( - None if clearedCategoryIds is None else frozenset(clearedCategoryIds) + None if clearedCategoryIds is None else clearedCategoryIds ) def plainText( diff --git a/Furious/Window/LogPage.py b/Furious/Window/LogPage.py index c3234bd9..cd75f0b3 100644 --- a/Furious/Window/LogPage.py +++ b/Furious/Window/LogPage.py @@ -344,6 +344,12 @@ class LogPage(Mixins.QTranslatable, QMainWindow): self.autoScrollSwitch.toggled.connect(self._autoScrollChanged) self.autoClearSwitch.toggled.connect(self._autoClearChanged) self.manager.categoryRegistered.connect(self._categoryRegistered) + # Do not drive the document directly from entryAdded: cross-thread Qt + # delivery can occur after that entry was evicted or its generation was + # cleared, and one queued signal per line defeats batching under a burst. + # entriesChanged coalesces producers, then one indexed snapshot provides + # an atomic truth that _synchronizeDocument applies as prefix removal and + # suffix append rather than rebuilding the document in steady state. self.manager.entriesChanged.connect(self._entriesChanged) self.manager.entriesCleared.connect(self._entriesCleared) diff --git a/tests/README.md b/tests/README.md index af6cbbd9..c8d3aff6 100644 --- a/tests/README.md +++ b/tests/README.md @@ -13,6 +13,7 @@ clients, or real proxy cores. | Area | Principal tests | | --- | --- | | Configuration, profiles, migration, repositories | `test_models_and_services.py`, `test_repository_contracts.py` | +| Generation log invariants, model fuzzing, concurrency, reclamation, complexity, and opt-in soak/latency probes | `test_log_manager_generation.py` | | Low-level application, runtime, editor, and storage contracts | `test_interface.py` | | Application composition, startup rollback, connection ownership, entry-point and crash boundaries | `test_architecture_refactors.py`, `test_application_process.py` | | Plugin registration, capability dispatch, factories, rollback, and Hysteria1 ownership | `test_plugin_architecture.py`, `test_hysteria1_protocol.py` | @@ -92,6 +93,9 @@ python -m unittest tests.test_qt_stress tests.test_process_stress -v # POSIX shell: export FURIOUS_VERY_HEAVY_TESTS=1 python -m unittest tests.test_very_heavy -v +# Generation-specific release-confidence campaign (same opt-in switch) +python -m unittest tests.test_log_manager_generation.VeryHeavyGenerationLogManagerTest -v + # Shared-state order-independence spot check python -m unittest tests.test_public_api tests.test_stylesheet_states tests.test_ui_behavior tests.test_dialog_geometry tests.test_main_window_geometry tests.test_isolation_and_navigation tests.test_frozenlib tests.test_service_runtime tests.test_endpoint_info tests.test_metrics_behavior tests.test_native_tun_semantics tests.test_backend_editor_contract tests.test_shadowsocks_uri tests.test_socks_uri tests.test_subscription_sync tests.test_subscription_manager tests.test_controllers tests.test_hysteria2_compatibility tests.test_hysteria1_protocol tests.test_plugin_architecture tests.test_architecture_refactors tests.test_repository_contracts tests.test_models_and_services tests.test_interface -v python -m unittest discover -s tests -v diff --git a/tests/test_log_manager_generation.py b/tests/test_log_manager_generation.py new file mode 100644 index 00000000..6149aca7 --- /dev/null +++ b/tests/test_log_manager_generation.py @@ -0,0 +1,1430 @@ +# Copyright (C) 2024–present Loren Eteval & contributors +# +# This file is part of Furious. +# +# This program is free software: you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation, either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . + +"""Adversarial contracts for the generation-based structured log store. + +The fast tests emphasize structural evidence and an independent reference +model. The opt-in class adds high-count timing, reclamation, and soak probes; +enable it with the repository-wide ``FURIOUS_VERY_HEAVY_TESTS`` switch. +""" + +from __future__ import annotations + +from Furious.Models import LogCategory +from Furious.Service.LogManager import ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + LogManager, +) +from Furious.Window.LogPage import LogPage + +from collections import OrderedDict +from datetime import datetime +from statistics import median + +import gc +import json +import os +import random +import subprocess +import threading +import time +import types +import unittest +import weakref + +from tests.support import ( + application, + collectAtBoundary, + isolatedSettings, + processQtEvents, + resourceSnapshot, + veryHeavyEnabled, + waitFor, +) + + +class _ObservedEntries(OrderedDict): + """Count container primitives without changing OrderedDict semantics.""" + + def __init__(self, entries=()): + self.iterations = 0 + self.deletions = 0 + self.oldestRemovals = 0 + super().__init__(entries) + + def __iter__(self): + self.iterations += 1 + return super().__iter__() + + def values(self): + self.iterations += 1 + return super().values() + + def items(self): + self.iterations += 1 + return super().items() + + def keys(self): + self.iterations += 1 + return super().keys() + + def __delitem__(self, key): + self.deletions += 1 + return super().__delitem__(key) + + def popitem(self, last=True): + if not last: + self.oldestRemovals += 1 + return super().popitem(last=last) + + +class _ReferenceLog: + """Slow, deliberately flat model of public LogManager semantics.""" + + def __init__(self, manager): + self.maximumEntries = manager.maximumEntries + self.maximumCharacters = manager.maximumCharacters + self.autoClearMaximumEntries = manager.autoClearMaximumEntries + self.autoClearEnabled = manager.autoClearEnabled + self.categories = {category.id: category for category in manager.categories()} + self.entries = [] + + def register(self, category): + self.categories.setdefault(category.id, category) + + def beforeAppend(self, categoryId): + if ( + categoryId == CORE_LOG_CATEGORY + and self.autoClearEnabled + and self.count(CORE_LOG_CATEGORY) >= self.autoClearMaximumEntries + ): + self.entries = [ + entry + for entry in self.entries + if entry.categoryId == APPLICATION_LOG_CATEGORY + ] + + def accept(self, entry): + self.entries.append(entry) + while self.entries and ( + len(self.entries) > self.maximumEntries + or sum(len(item.message) for item in self.entries) > self.maximumCharacters + ): + del self.entries[0] + + def clear(self, categoryId=None, *, runtimeOnly=False): + if categoryId is None: + if runtimeOnly: + self.entries = [ + entry + for entry in self.entries + if not self.categories[entry.categoryId].runtime + ] + else: + self.entries.clear() + return + + self.entries = [ + entry for entry in self.entries if entry.categoryId != categoryId + ] + + def setAutoClearEnabled(self, enabled): + self.autoClearEnabled = bool(enabled) + if ( + self.autoClearEnabled + and self.count(CORE_LOG_CATEGORY) >= self.autoClearMaximumEntries + ): + self.entries = [ + entry + for entry in self.entries + if entry.categoryId == APPLICATION_LOG_CATEGORY + ] + + def count(self, categoryId): + return sum(entry.categoryId == categoryId for entry in self.entries) + + +def _batchEntries(batch): + """Enumerate one retired owner without trusting manager counters.""" + return tuple(batch.entries.values()) + + +def _assertManagerInvariants(testCase, manager, model=None): + """Independently prove ownership, indexing, accounting, and ordering. + + This helper intentionally derives truth from the actual entry objects in + both ownership indexes. It then compares the public API and every cached + aggregate against that derivation, so a defect cannot hide merely because + two manager counters drift in the same direction. + """ + with manager._lock: + generations = manager._activeGenerationsLocked() + testCase.assertEqual(len(generations), 3) + testCase.assertEqual(len({id(item) for item in generations}), 3) + testCase.assertEqual(len({item.identifier for item in generations}), 3) + testCase.assertEqual( + tuple(item.scope for item in generations), + ('application', 'runtime', 'other'), + ) + + liveEntries = [] + liveObjectIds = set() + categoryTruth = {category.id: [] for category in manager.categories()} + + for generation in generations: + chronological = tuple(generation.entries.values()) + testCase.assertEqual( + tuple(entry.sequence for entry in chronological), + tuple(sorted(entry.sequence for entry in chronological)), + ) + testCase.assertEqual( + generation.characterCount, + sum(len(entry.message) for entry in chronological), + ) + + indexedSequences = [] + for categoryId, categoryEntries in generation.entriesByCategory.items(): + category = manager._categories[categoryId] + testCase.assertIs( + manager._generationForCategoryLocked(category), generation + ) + indexed = tuple(categoryEntries.entries.values()) + testCase.assertTrue(indexed) + testCase.assertEqual( + tuple(entry.sequence for entry in indexed), + tuple(sorted(entry.sequence for entry in indexed)), + ) + testCase.assertTrue( + all(entry.categoryId == categoryId for entry in indexed) + ) + testCase.assertEqual( + categoryEntries.characterCount, + sum(len(entry.message) for entry in indexed), + ) + indexedSequences.extend(entry.sequence for entry in indexed) + categoryTruth[categoryId].extend(indexed) + + testCase.assertCountEqual( + indexedSequences, + (entry.sequence for entry in chronological), + ) + for entry in chronological: + testCase.assertNotIn(id(entry), liveObjectIds) + liveObjectIds.add(id(entry)) + liveEntries.extend(chronological) + + liveEntries.sort(key=lambda entry: entry.sequence) + testCase.assertEqual( + tuple(entry.sequence for entry in liveEntries), + tuple(sorted({entry.sequence for entry in liveEntries})), + ) + testCase.assertEqual(manager._retainedEntryCount, len(liveEntries)) + testCase.assertEqual( + manager._retainedCharacters, + sum(len(entry.message) for entry in liveEntries), + ) + testCase.assertGreaterEqual(manager._retainedEntryCount, 0) + testCase.assertGreaterEqual(manager._retainedCharacters, 0) + + retiredEntries = [] + for batch in manager._retiredBatches: + entries = _batchEntries(batch) + testCase.assertEqual( + batch.characterCount, + sum(len(entry.message) for entry in entries), + ) + retiredEntries.extend(entries) + + retiredIds = {id(entry) for entry in retiredEntries} + testCase.assertTrue(liveObjectIds.isdisjoint(retiredIds)) + testCase.assertEqual(len(retiredEntries), manager._retiredEntryCount) + testCase.assertEqual( + sum(len(entry.message) for entry in retiredEntries), + manager._retiredCharacters, + ) + testCase.assertGreaterEqual(manager._retiredEntryCount, 0) + testCase.assertGreaterEqual(manager._retiredCharacters, 0) + + publicEntries = manager.entries() + testCase.assertEqual(publicEntries, tuple(liveEntries)) + testCase.assertTrue(retiredIds.isdisjoint(id(entry) for entry in publicEntries)) + testCase.assertEqual(manager.entryCount(), len(liveEntries)) + testCase.assertEqual( + manager.retainedCharacters, + sum(len(entry.message) for entry in liveEntries), + ) + + for category in manager.categories(): + expected = tuple(categoryTruth[category.id]) + testCase.assertEqual(manager.entries(category.id), expected) + testCase.assertEqual(manager.entryCount(category.id), len(expected)) + + if model is not None: + testCase.assertEqual( + tuple(entry.sequence for entry in manager.entries()), + tuple(entry.sequence for entry in model.entries), + ) + testCase.assertEqual( + tuple(entry.message for entry in manager.entries()), + tuple(entry.message for entry in model.entries), + ) + + +def _percentile(samples, percentage): + """Return a nearest-rank percentile from nanosecond samples.""" + ordered = sorted(samples) + index = min(len(ordered) - 1, int((len(ordered) - 1) * percentage)) + return ordered[index] / 1_000 + + +class GenerationLogManagerContractTest(unittest.TestCase): + """Exercise exact semantics and structural complexity in the fast tier.""" + + @classmethod + def setUpClass(cls): + application() + + def makeManager(self, **kwargs): + defaults = { + 'maximumEntries': 37, + 'maximumCharacters': 400, + 'maximumEntryCharacters': 80, + 'autoClearMaximumEntries': 5, + } + defaults.update(kwargs) + manager = LogManager(**defaults) + manager.registerComponent('runtime.extra', 'Runtime extra', runtime=True) + manager.registerComponent('other.extra', 'Other extra', runtime=False) + return manager + + def appendBoth(self, manager, model, message, categoryId): + model.beforeAppend(categoryId) + entry = manager.append(message, categoryId) + model.accept(entry) + _assertManagerInvariants(self, manager, model) + return entry + + def testDeterministicStateTransitionMatrix(self): + """Mix every clear kind, rollover, retention mode, and stream.""" + manager = self.makeManager( + maximumEntries=7, + maximumCharacters=35, + maximumEntryCharacters=35, + ) + model = _ReferenceLog(manager) + operations = ( + ('a', APPLICATION_LOG_CATEGORY), + ('r', 'runtime.extra'), + ('o', 'other.extra'), + ('c1', CORE_LOG_CATEGORY), + ('t', TUN2SOCKS_LOG_CATEGORY), + ('c2', CORE_LOG_CATEGORY), + ) + for message, categoryId in operations: + self.appendBoth(manager, model, message, categoryId) + + manager.clear('runtime.extra') + model.clear('runtime.extra') + _assertManagerInvariants(self, manager, model) + manager.clear(runtimeOnly=True) + model.clear(runtimeOnly=True) + _assertManagerInvariants(self, manager, model) + + for index in range(9): + self.appendBoth(manager, model, f'core-{index}', CORE_LOG_CATEGORY) + + manager.setAutoClearEnabled(False) + model.setAutoClearEnabled(False) + manager.clear('other.extra') + model.clear('other.extra') + self.appendBoth(manager, model, 'x' * 80, 'other.extra') + self.appendBoth(manager, model, 'y' * 81, APPLICATION_LOG_CATEGORY) + manager.clear() + model.clear() + _assertManagerInvariants(self, manager, model) + self.appendBoth(manager, model, 'after-clear', APPLICATION_LOG_CATEGORY) + + def testSeededModelBasedStateMachine(self): + """Compare arbitrary public transitions to a flat reference model.""" + seeds = (0, 1, 7, 19, 41, 97, 313, 997) + for seed in seeds: + with self.subTest(seed=seed): + randomizer = random.Random(seed) + manager = self.makeManager() + manager.RetiredCleanupBudget = randomizer.choice((1, 2, 64)) + model = _ReferenceLog(manager) + categories = tuple(model.categories) + history = [] + + for operationIndex in range(350): + operation = randomizer.randrange(100) + try: + if operation < 66: + categoryId = randomizer.choice(categories) + message = randomizer.choice( + ('', 'x', 'line\n', '😀é', '\0', 'z' * 93) + ) + history.append(('append', categoryId, len(message))) + self.appendBoth(manager, model, message, categoryId) + elif operation < 73: + categoryId = randomizer.choice(categories) + history.append(('clear-category', categoryId)) + manager.clear(categoryId) + model.clear(categoryId) + elif operation < 80: + history.append(('clear-runtime',)) + manager.clear(runtimeOnly=True) + model.clear(runtimeOnly=True) + elif operation < 84: + history.append(('clear-all',)) + manager.clear() + model.clear() + elif operation < 91: + enabled = bool(randomizer.getrandbits(1)) + history.append(('auto-clear', enabled)) + manager.setAutoClearEnabled(enabled) + model.setAutoClearEnabled(enabled) + else: + history.append(('read',)) + manager.snapshot(randomizer.choice((None, *categories))) + + _assertManagerInvariants(self, manager, model) + except Exception as error: + self.fail( + f'seed={seed} operation={operationIndex} error={error!r} ' + f'history={history!r}' + ) + + def testRolloverAndClearUseConstantStructuralWork(self): + """Prove logical swaps never traverse, unlink, or pop old entries.""" + for size in (1, 100, 1_000, 10_000): + with self.subTest(size=size): + manager = self.makeManager( + maximumEntries=size * 3 + 10, + maximumCharacters=size * 30 + 100, + autoClearMaximumEntries=size, + ) + for index in range(size): + manager.append(f'core {index}', CORE_LOG_CATEGORY) + manager.append(f'other {index}', 'other.extra') + + observed = [] + for generation in manager._activeGenerationsLocked(): + index = _ObservedEntries(generation.entries) + generation.entries = index + observed.append(index) + + manager.append('trigger', CORE_LOG_CATEGORY) + self.assertEqual(sum(item.iterations for item in observed), 0) + self.assertEqual(sum(item.deletions for item in observed), 0) + self.assertEqual(sum(item.oldestRemovals for item in observed), 0) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 1) + self.assertEqual(manager.entryCount('other.extra'), 0) + + for index in range(size): + manager.append(f'runtime {index}', 'runtime.extra') + runtime = manager._runtimeGeneration + watchedRuntime = _ObservedEntries(runtime.entries) + runtime.entries = watchedRuntime + manager.clear(runtimeOnly=True) + self.assertEqual(watchedRuntime.iterations, 0) + self.assertEqual(watchedRuntime.oldestRemovals, 0) + + manager.append('application', APPLICATION_LOG_CATEGORY) + watched = [] + for generation in manager._activeGenerationsLocked(): + index = _ObservedEntries(generation.entries) + generation.entries = index + watched.append(index) + manager.clear() + self.assertEqual(sum(item.iterations for item in watched), 0) + self.assertEqual(sum(item.oldestRemovals for item in watched), 0) + _assertManagerInvariants(self, manager) + + def testEmptyRolloversDoNotQueueBatchesOrReuseGenerationIds(self): + """Swap empty streams repeatedly without retired-queue amplification.""" + manager = self.makeManager(autoClearEnabled=False) + identifiers = set() + for _index in range(2_000): + manager.clear() + current = tuple( + generation.identifier + for generation in manager._activeGenerationsLocked() + ) + self.assertTrue(identifiers.isdisjoint(current)) + identifiers.update(current) + self.assertEqual(manager.retiredEntryCount, 0) + self.assertEqual(len(manager._retiredBatches), 0) + _assertManagerInvariants(self, manager) + + def testSnapshotAndOldestEvictionInspectOnlyFixedIndexes(self): + """Count one category traversal and at most three live stream heads.""" + manager = self.makeManager( + maximumEntries=3, + maximumCharacters=100, + maximumEntryCharacters=100, + autoClearEnabled=False, + ) + for categoryId in ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + 'other.extra', + ): + manager.append(categoryId, categoryId) + + globalIndexes = [] + categoryIndexes = [] + for generation in manager._activeGenerationsLocked(): + observed = _ObservedEntries(generation.entries) + generation.entries = observed + globalIndexes.append(observed) + for categoryEntries in generation.entriesByCategory.values(): + observed = _ObservedEntries(categoryEntries.entries) + categoryEntries.entries = observed + categoryIndexes.append(observed) + + manager.snapshot(CORE_LOG_CATEGORY) + self.assertEqual(sum(index.iterations for index in globalIndexes), 0) + self.assertEqual(sum(index.iterations for index in categoryIndexes), 1) + for index in (*globalIndexes, *categoryIndexes): + index.iterations = 0 + + manager.snapshot() + self.assertEqual(sum(index.iterations for index in globalIndexes), 3) + self.assertEqual(sum(index.iterations for index in categoryIndexes), 0) + for index in globalIndexes: + index.iterations = 0 + + manager.append('evict', CORE_LOG_CATEGORY) + self.assertLessEqual(sum(index.iterations for index in globalIndexes), 3) + self.assertEqual(sum(index.oldestRemovals for index in globalIndexes), 1) + _assertManagerInvariants(self, manager) + + def testCleanupBudgetIsGlobalFifoAndHandlesBoundaries(self): + """Limit one invocation across all batches, resuming the FIFO head.""" + for backlog, budget in ((1, 64), (63, 64), (64, 64), (65, 64), (130, 64)): + with self.subTest(backlog=backlog, budget=budget): + manager = self.makeManager(maximumEntries=200) + manager.RetiredCleanupBudget = budget + for index in range(backlog): + manager.append(str(index), CORE_LOG_CATEGORY) + if index in (1, 4): + manager.clear(runtimeOnly=True) + manager.clear(runtimeOnly=True) + before = manager.retiredEntryCount + identifiers = tuple( + getattr(batch, 'identifier', None) + for batch in manager._retiredBatches + ) + with manager._lock: + released = manager._cleanupRetiredLocked() + self.assertEqual(released, min(budget, before)) + self.assertEqual(manager.retiredEntryCount, before - released) + if manager._retiredBatches and identifiers: + self.assertIn( + getattr(manager._retiredBatches[0], 'identifier', None), + identifiers, + ) + _assertManagerInvariants(self, manager) + + manager = self.makeManager() + manager.append('retired', CORE_LOG_CATEGORY) + manager.clear(runtimeOnly=True) + with manager._lock: + self.assertEqual(manager._cleanupRetiredLocked(0), 0) + manager.RetiredCleanupBudget = 0 + with manager._lock: + self.assertEqual(manager._cleanupRetiredLocked(), 1) + + def testRetiredGenerationCleanupIsStrictlyFifo(self): + """Finish each older clear-all stream before touching the next one.""" + manager = self.makeManager(autoClearEnabled=False) + manager.RetiredCleanupBudget = 1 + manager.append('application 1', APPLICATION_LOG_CATEGORY) + manager.append('application 2', APPLICATION_LOG_CATEGORY) + manager.append('runtime 1', CORE_LOG_CATEGORY) + manager.append('runtime 2', CORE_LOG_CATEGORY) + manager.append('other 1', 'other.extra') + manager.append('other 2', 'other.extra') + expectedIdentifiers = tuple( + generation.identifier for generation in manager._activeGenerationsLocked() + ) + manager.clear() + self.assertEqual( + tuple(batch.identifier for batch in manager._retiredBatches), + expectedIdentifiers, + ) + observedHeads = [] + while manager.retiredEntryCount: + observedHeads.append(manager._retiredBatches[0].identifier) + with manager._lock: + self.assertEqual(manager._cleanupRetiredLocked(), 1) + self.assertEqual( + observedHeads, + [ + expectedIdentifiers[0], + expectedIdentifiers[0], + expectedIdentifiers[1], + expectedIdentifiers[1], + expectedIdentifiers[2], + expectedIdentifiers[2], + ], + ) + + def testWeakReferencesSnapshotsAndExternalOwners(self): + """Separate manager ownership from snapshot and caller ownership.""" + for size in (1, 63, 64, 65, 1_000, 10_000): + with self.subTest(size=size): + manager = self.makeManager( + maximumEntries=size + 2, + maximumCharacters=(size + 2) * 16, + maximumEntryCharacters=16, + autoClearEnabled=False, + ) + manager.RetiredCleanupBudget = 64 + for index in range(size): + manager.append(f'entry {index}', CORE_LOG_CATEGORY) + historical = manager.entries(CORE_LOG_CATEGORY) + self.assertEqual(len(historical), size) + references = tuple(weakref.ref(entry) for entry in historical) + externallyOwned = historical[-1] + manager.clear(runtimeOnly=True) + self.assertTrue( + all(reference() is not None for reference in references) + ) + while manager.retiredEntryCount: + with manager._lock: + manager._cleanupRetiredLocked() + self.assertTrue( + all(reference() is not None for reference in references) + ) + del historical + gc.collect() + self.assertTrue( + all(reference() is None for reference in references[:-1]) + ) + self.assertIs(references[-1](), externallyOwned) + del externallyOwned + gc.collect() + self.assertIsNone(references[-1]()) + + def testRetentionThreeWayMergeAndLargeSequences(self): + """Evict only the global oldest live head across all active streams.""" + manager = self.makeManager( + maximumEntries=5, + maximumCharacters=19, + maximumEntryCharacters=10, + autoClearEnabled=False, + ) + manager._sequence = 10**40 + categories = ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + 'other.extra', + APPLICATION_LOG_CATEGORY, + 'runtime.extra', + 'other.extra', + CORE_LOG_CATEGORY, + ) + expected = [] + for index, categoryId in enumerate(categories): + entry = manager.append(str(index) * (index % 4 + 1), categoryId) + expected.append(entry) + while ( + len(expected) > manager.maximumEntries + or sum(len(item.message) for item in expected) + > manager.maximumCharacters + ): + del expected[0] + self.assertEqual(manager.entries(), tuple(expected)) + _assertManagerInvariants(self, manager) + self.assertTrue(all(entry.sequence > 10**40 for entry in manager.entries())) + + def testSelectiveClearAndRegistrationStress(self): + """Keep shared stream indexes exact for tiny and dominant categories.""" + manager = self.makeManager( + maximumEntries=5_000, + maximumCharacters=100_000, + autoClearEnabled=False, + ) + categoryIds = [] + for index in range(100): + category = manager.registerCategory( + LogCategory( + f'dynamic.{index}', + f'Dynamic {index}', + runtime=bool(index % 2), + translatable=bool(index % 3), + ) + ) + categoryIds.append(category.id) + for index in range(2_000): + manager.append(str(index), categoryIds[index % len(categoryIds)]) + for categoryId in (categoryIds[0], categoryIds[-1], categoryIds[51]): + expected = tuple( + entry for entry in manager.entries() if entry.categoryId != categoryId + ) + manager.clear(categoryId) + self.assertEqual(manager.entries(), expected) + _assertManagerInvariants(self, manager) + duplicate = manager.category(categoryIds[1]) + self.assertIs(manager.registerCategory(duplicate), duplicate) + with self.assertRaises(ValueError): + manager.registerCategory(LogCategory(categoryIds[1], 'Different')) + + def testPathologicalMessagesAndRejectedOperationsAreAtomic(self): + """Account normalized storage and preserve state after bad input.""" + manager = self.makeManager(maximumEntryCharacters=32) + messages = ('', 'x', 'line\r\n', '😀é', 'a\0b', 'n\n' * 20, 'z' * 33) + for message in messages: + entry = manager.append(message) + self.assertLessEqual(len(entry.message), 32) + _assertManagerInvariants(self, manager) + + class RaisingString: + def __str__(self): + raise RuntimeError('conversion failed') + + before = manager.entries() + with self.assertRaises(RuntimeError): + manager.append(RaisingString()) + with self.assertRaises(TypeError): + manager.append('bad timestamp', timestamp='not a datetime') + with self.assertRaises(KeyError): + manager.append('unknown', 'missing') + with self.assertRaises(ValueError): + manager.clear(CORE_LOG_CATEGORY, runtimeOnly=True) + self.assertEqual(manager.entries(), before) + _assertManagerInvariants(self, manager) + + def testSignalsAreOutsideLockReentrantAndCoalesced(self): + """Allow querying/clearing slots without deadlock or stale UI truth.""" + manager = self.makeManager(autoClearEnabled=False) + added = [] + cleared = [] + changed = [] + lockWasFreeDuringSignal = [] + + def onAdded(entry): + added.append(entry.sequence) + manager.snapshot() + manager.entryCount() + if not lockWasFreeDuringSignal: + completed = threading.Event() + + def queryFromAnotherThread(): + manager.entryCount() + completed.set() + + worker = threading.Thread(target=queryFromAnotherThread) + worker.start() + lockWasFreeDuringSignal.append(completed.wait(2)) + worker.join(2) + if entry.categoryId == CORE_LOG_CATEGORY: + manager.clear(runtimeOnly=True) + + manager.entryAdded.connect(onAdded) + manager.entriesCleared.connect(cleared.append) + manager.entriesChanged.connect(changed.append) + manager.append('core', CORE_LOG_CATEGORY) + for index in range(20): + manager.append(f'application {index}') + processQtEvents() + self.assertEqual(len(added), 21) + self.assertEqual(cleared, [manager._runtimeCategoryIds]) + self.assertEqual(changed, [21]) + self.assertEqual(lockWasFreeDuringSignal, [True]) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 0) + _assertManagerInvariants(self, manager) + + def testAppendAndClearCannotSplitOneAtomicMutation(self): + """Force clear to wait while append has selected its generation.""" + manager = self.makeManager(autoClearEnabled=False) + selected = threading.Event() + release = threading.Event() + appendThreadId = [] + original = manager._generationForCategoryLocked + + def observedGeneration(category): + generation = original(category) + if threading.get_ident() in appendThreadId and not selected.is_set(): + selected.set() + self.assertTrue(release.wait(5)) + return generation + + manager._generationForCategoryLocked = observedGeneration + errors = [] + + def append(): + appendThreadId.append(threading.get_ident()) + try: + manager.append('racing', CORE_LOG_CATEGORY) + except Exception as error: + errors.append(error) + + producer = threading.Thread(target=append) + clearer = threading.Thread(target=lambda: manager.clear(runtimeOnly=True)) + producer.start() + self.assertTrue(selected.wait(5)) + clearer.start() + time.sleep(0.01) + self.assertTrue(clearer.is_alive()) + release.set() + producer.join(5) + clearer.join(5) + self.assertFalse(producer.is_alive()) + self.assertFalse(clearer.is_alive()) + self.assertEqual(errors, []) + self.assertEqual(manager.entries(CORE_LOG_CATEGORY), tuple()) + _assertManagerInvariants(self, manager) + + def testConcurrentReadersWritersAndMutatorsFinishConsistently(self): + """Stress every supported locked operation with deterministic seeds.""" + for producerCount in (1, 2, 4, 8, 16): + with self.subTest(producers=producerCount): + manager = self.makeManager( + maximumEntries=500, + maximumCharacters=20_000, + autoClearMaximumEntries=17, + ) + errors = [] + returnedSequences = [] + sequenceLock = threading.Lock() + # Producers plus the reader and mutator are the complete + # participant set. Keeping the exact cardinality here makes a + # failed synchronization a real deadlock signal rather than a + # test harness waiting for a thread that was never created. + start = threading.Barrier(producerCount + 2) + + def producer(workerIndex): + randomizer = random.Random(10_000 + workerIndex) + try: + start.wait(5) + for index in range(300): + categoryId = randomizer.choice( + ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + 'runtime.extra', + 'other.extra', + ) + ) + entry = manager.append(f'{workerIndex}:{index}', categoryId) + with sequenceLock: + returnedSequences.append(entry.sequence) + except Exception as error: + errors.append(error) + + def reader(): + try: + start.wait(5) + for _index in range(500): + entries = manager.entries() + self.assertEqual( + tuple(item.sequence for item in entries), + tuple(sorted(item.sequence for item in entries)), + ) + manager.entryCount(CORE_LOG_CATEGORY) + except Exception as error: + errors.append(error) + + def mutator(): + randomizer = random.Random(77) + try: + start.wait(5) + for index in range(180): + choice = randomizer.randrange(4) + if choice == 0: + manager.clear(runtimeOnly=True) + elif choice == 1: + manager.clear('other.extra') + elif choice == 2: + manager.setAutoClearEnabled(bool(index % 2)) + else: + manager.snapshot(CORE_LOG_CATEGORY) + except Exception as error: + errors.append(error) + + threads = [ + threading.Thread(target=producer, args=(index,)) + for index in range(producerCount) + ] + threads.extend( + (threading.Thread(target=reader), threading.Thread(target=mutator)) + ) + for thread in threads: + thread.start() + for thread in threads: + thread.join(20) + self.assertFalse(thread.is_alive(), 'possible LogManager deadlock') + self.assertEqual(errors, []) + self.assertEqual(len(returnedSequences), producerCount * 300) + self.assertEqual(len(set(returnedSequences)), len(returnedSequences)) + _assertManagerInvariants(self, manager) + + def testLogPageMatchesTruthAcrossPruneRolloverAndFilters(self): + """Keep the incremental document exact through generation mutations.""" + with isolatedSettings(): + manager = self.makeManager( + maximumEntries=120, + maximumCharacters=10_000, + autoClearMaximumEntries=40, + ) + page = LogPage(manager=manager) + page.resize(900, 420) + page.show() + for index in range(400): + categoryId = ( + APPLICATION_LOG_CATEGORY if index % 7 == 0 else CORE_LOG_CATEGORY + ) + manager.append(f'line {index:04d}', categoryId) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertEqual( + page.plainText().splitlines(), + [entry.message for entry in manager.entries()], + ) + for categoryId in (CORE_LOG_CATEGORY, APPLICATION_LOG_CATEGORY, 'all'): + page.filterComboBox.setCurrentIndex( + page.filterComboBox.findData(categoryId) + ) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertEqual( + page.plainText().splitlines(), + [entry.message for entry in manager.entries(categoryId)], + ) + manager.clear(runtimeOnly=True) + page.filterComboBox.setCurrentIndex(page.filterComboBox.findData('all')) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertEqual( + page.plainText().splitlines(), + [entry.message for entry in manager.entries()], + ) + page.close() + page.deleteLater() + collectAtBoundary() + + +@unittest.skipUnless( + veryHeavyEnabled(), + 'set FURIOUS_VERY_HEAVY_TESTS=1 for generation stress/benchmarks', +) +class VeryHeavyGenerationLogManagerTest(unittest.TestCase): + """Run release-confidence scaling, latency, backlog, and soak probes.""" + + @classmethod + def setUpClass(cls): + application() + + def report(self, name, **values): + print( + 'GENERATION_CAMPAIGN_REPORT=' + + json.dumps({'name': name, **values}, sort_keys=True), + flush=True, + ) + + def testScalingAndStructuralComplexity(self): + """Measure geometric scaling while structural assertions guard O(1).""" + results = [] + for size in (100, 1_000, 10_000, 100_000): + manager = LogManager( + maximumEntries=size + 10, + maximumCharacters=(size + 10) * 16, + maximumEntryCharacters=16, + autoClearMaximumEntries=size, + ) + for index in range(size): + manager.append(str(index), CORE_LOG_CATEGORY) + samples = [] + synchronousTraversals = 0 + synchronousRemovals = 0 + for repetition in range(7): + if repetition: + for index in range(size): + manager.append(str(index), CORE_LOG_CATEGORY) + watched = _ObservedEntries(manager._runtimeGeneration.entries) + manager._runtimeGeneration.entries = watched + started = time.perf_counter_ns() + manager.clear(runtimeOnly=True) + samples.append(time.perf_counter_ns() - started) + synchronousTraversals += watched.iterations + synchronousRemovals += watched.oldestRemovals + self.assertEqual(watched.iterations, 0) + self.assertEqual(watched.oldestRemovals, 0) + results.append( + { + 'n': size, + 'runtime_clear_us': median(samples) / 1_000, + 'synchronous_traversals': synchronousTraversals, + 'synchronous_removals': synchronousRemovals, + } + ) + ratio = results[-1]['runtime_clear_us'] / max( + results[0]['runtime_clear_us'], 0.001 + ) + self.assertLess(ratio, 100) + self.report('scaling', samples=results, endpoint_ratio=ratio) + + def testOperationScalingMatrix(self): + """Report medians for every claimed constant, linear, or bounded path.""" + + def elapsed(operation): + started = time.perf_counter_ns() + operation() + return time.perf_counter_ns() - started + + matrix = [] + for size in (100, 1_000, 10_000): + samples = { + name: [] + for name in ( + 'normal_append', + 'core_rollover', + 'runtime_clear', + 'clear_all', + 'sole_category_clear', + 'shared_category_clear', + 'cleanup_64', + 'category_snapshot', + 'global_snapshot', + 'oldest_eviction', + ) + } + for _repetition in range(5): + manager = LogManager( + maximumEntries=size * 2 + 10, + autoClearMaximumEntries=size, + ) + for index in range(size): + manager.append(str(index), CORE_LOG_CATEGORY) + samples['normal_append'].append( + elapsed(lambda: manager.append('ordinary')) + ) + samples['core_rollover'].append( + elapsed(lambda: manager.append('trigger', CORE_LOG_CATEGORY)) + ) + + manager = LogManager(maximumEntries=size + 10, autoClearEnabled=False) + for index in range(size): + manager.append(str(index), CORE_LOG_CATEGORY) + samples['category_snapshot'].append( + elapsed(lambda: manager.snapshot(CORE_LOG_CATEGORY)) + ) + samples['runtime_clear'].append( + elapsed(lambda: manager.clear(runtimeOnly=True)) + ) + with manager._lock: + samples['cleanup_64'].append( + elapsed(lambda: manager._cleanupRetiredLocked()) + ) + + manager = LogManager(maximumEntries=size + 10, autoClearEnabled=False) + for index in range(size): + manager.append(str(index), APPLICATION_LOG_CATEGORY) + samples['sole_category_clear'].append( + elapsed(lambda: manager.clear(APPLICATION_LOG_CATEGORY)) + ) + + manager = LogManager( + maximumEntries=size * 2 + 10, + autoClearEnabled=False, + ) + for index in range(size): + manager.append(str(index), CORE_LOG_CATEGORY) + manager.append('shared', TUN2SOCKS_LOG_CATEGORY) + samples['shared_category_clear'].append( + elapsed(lambda: manager.clear(CORE_LOG_CATEGORY)) + ) + + manager = LogManager(maximumEntries=size + 10, autoClearEnabled=False) + for index in range(size): + manager.append( + str(index), + ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + )[index % 3], + ) + samples['global_snapshot'].append(elapsed(lambda: manager.snapshot())) + samples['clear_all'].append(elapsed(lambda: manager.clear())) + + manager = LogManager(maximumEntries=size, autoClearEnabled=False) + for index in range(size): + manager.append( + str(index), + ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + )[index % 3], + ) + samples['oldest_eviction'].append( + elapsed(lambda: manager.append('evict')) + ) + + matrix.append( + { + 'n': size, + **{ + name + '_us': median(values) / 1_000 + for name, values in samples.items() + }, + } + ) + + ratios = { + name: matrix[-1][name] / max(matrix[0][name], 0.001) + for name in matrix[0] + if name != 'n' + } + for name in ( + 'normal_append_us', + 'core_rollover_us', + 'runtime_clear_us', + 'clear_all_us', + 'sole_category_clear_us', + 'cleanup_64_us', + 'oldest_eviction_us', + ): + self.assertLess(ratios[name], 50) + self.report('operation-scaling-matrix', samples=matrix, ratios=ratios) + + def testAppendLatencyWithRolloverRetentionAndBacklog(self): + """Report latency distribution by append-path condition.""" + manager = LogManager( + maximumEntries=5_000, + maximumCharacters=200_000, + maximumEntryCharacters=128, + autoClearMaximumEntries=100, + ) + manager.RetiredCleanupBudget = 64 + ordinary = [] + rollover = [] + retention = [] + maximumRetired = 0 + maximumBatches = 0 + for index in range(50_000): + categoryId = CORE_LOG_CATEGORY if index % 3 else APPLICATION_LOG_CATEGORY + triggersRollover = ( + categoryId == CORE_LOG_CATEGORY + and manager.entryCount(CORE_LOG_CATEGORY) + >= manager.autoClearMaximumEntries + ) + atRetention = manager.entryCount() >= manager.maximumEntries + message = f'{index:06d}-' + 'x' * (index % 97) + causesRetention = ( + manager.entryCount() + 1 > manager.maximumEntries + or manager.retainedCharacters + len(message) > manager.maximumCharacters + ) + started = time.perf_counter_ns() + manager.append(message, categoryId) + elapsed = time.perf_counter_ns() - started + if triggersRollover: + rollover.append(elapsed) + elif atRetention or causesRetention: + retention.append(elapsed) + else: + ordinary.append(elapsed) + maximumRetired = max(maximumRetired, manager.retiredEntryCount) + maximumBatches = max(maximumBatches, len(manager._retiredBatches)) + + def distribution(samples): + return { + 'count': len(samples), + 'median_us': median(samples) / 1_000, + 'p95_us': _percentile(samples, 0.95), + 'p99_us': _percentile(samples, 0.99), + 'p999_us': _percentile(samples, 0.999), + 'max_us': max(samples) / 1_000, + } + + self.assertTrue(ordinary) + self.assertTrue(rollover) + self.assertLessEqual(maximumRetired, manager.maximumEntries) + _assertManagerInvariants(self, manager) + self.report( + 'append-latency', + ordinary=distribution(ordinary), + rollover=distribution(rollover), + retention=distribution(retention) if retention else None, + maximum_retired=maximumRetired, + maximum_retired_batches=maximumBatches, + ) + + def testCleanupLatencyDoesNotScaleWithBacklog(self): + """Keep append cleanup capped at 64 for geometrically larger queues.""" + results = [] + for backlog in (64, 1_000, 10_000): + samples = [] + for _repetition in range(9): + manager = LogManager( + maximumEntries=backlog + 10, + maximumCharacters=(backlog + 10) * 8, + maximumEntryCharacters=8, + autoClearEnabled=False, + ) + manager.RetiredCleanupBudget = 64 + for index in range(backlog): + manager.append(str(index), CORE_LOG_CATEGORY) + manager.clear(runtimeOnly=True) + before = manager.retiredEntryCount + started = time.perf_counter_ns() + manager.append('new') + samples.append(time.perf_counter_ns() - started) + self.assertEqual(manager.retiredEntryCount, max(0, before - 64)) + results.append({'backlog': backlog, 'median_us': median(samples) / 1_000}) + ratio = results[-1]['median_us'] / max(results[0]['median_us'], 0.001) + self.assertLess(ratio, 20) + self.report('cleanup-backlog-latency', samples=results, endpoint_ratio=ratio) + + def testAdversarialClearRateCannotGrowPhysicalBacklog(self): + """Retire generations as fast as the minimum cleanup rate permits.""" + capacity = 2_000 + manager = LogManager(maximumEntries=capacity, autoClearEnabled=False) + manager.RetiredCleanupBudget = 1 + maximumRetired = 0 + maximumPhysical = 0 + maximumBatches = 0 + for cycle in range(100): + for index in range(capacity): + manager.append(f'{cycle}:{index}', CORE_LOG_CATEGORY) + maximumRetired = max(maximumRetired, manager.retiredEntryCount) + maximumPhysical = max( + maximumPhysical, + manager.entryCount() + manager.retiredEntryCount, + ) + manager.clear(runtimeOnly=True) + maximumRetired = max(maximumRetired, manager.retiredEntryCount) + maximumPhysical = max( + maximumPhysical, + manager.entryCount() + manager.retiredEntryCount, + ) + maximumBatches = max(maximumBatches, len(manager._retiredBatches)) + self.assertLessEqual(maximumPhysical, capacity) + while manager.retiredEntryCount: + with manager._lock: + manager._cleanupRetiredLocked() + self.assertEqual(manager.retiredCharacters, 0) + self.assertEqual(len(manager._retiredBatches), 0) + self.report( + 'adversarial-backlog', + operations=capacity * 100, + cleanup_budget=1, + maximum_retired=maximumRetired, + maximum_physical_entries=maximumPhysical, + maximum_retired_batches=maximumBatches, + final_retired=manager.retiredEntryCount, + ) + + def testLongRunningModelSoakTracksMemoryAndBacklog(self): + """Sustain mixed operations while checking invariants and RSS plateaus.""" + seed = 0xF017105 + randomizer = random.Random(seed) + manager = LogManager( + maximumEntries=2_000, + maximumCharacters=200_000, + maximumEntryCharacters=256, + autoClearMaximumEntries=400, + ) + manager.RetiredCleanupBudget = 64 + manager.registerComponent('soak.runtime', 'Soak runtime', runtime=True) + manager.registerComponent('soak.other', 'Soak other', runtime=False) + categories = ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + 'soak.runtime', + 'soak.other', + ) + rssSamples = [] + maximumRetired = 0 + maximumRetiredCharacters = 0 + maximumBatches = 0 + maximumLive = 0 + maximumPhysical = 0 + started = time.perf_counter() + for index in range(250_000): + operation = randomizer.randrange(100) + if operation < 70: + manager.append( + f'{index}:{randomizer.randrange(10**9)}', + randomizer.choice(categories), + ) + elif operation < 80: + manager.snapshot(randomizer.choice((None, *categories))) + elif operation < 85: + manager.clear(randomizer.choice(categories)) + elif operation < 90: + manager.clear(runtimeOnly=True) + elif operation < 95: + manager.setAutoClearEnabled(bool(randomizer.getrandbits(1))) + else: + manager.entryCount(randomizer.choice((None, *categories))) + + maximumRetired = max(maximumRetired, manager.retiredEntryCount) + maximumRetiredCharacters = max( + maximumRetiredCharacters, manager.retiredCharacters + ) + maximumBatches = max(maximumBatches, len(manager._retiredBatches)) + maximumLive = max(maximumLive, manager.entryCount()) + maximumPhysical = max( + maximumPhysical, + manager.entryCount() + manager.retiredEntryCount, + ) + if (index + 1) % 10_000 == 0: + _assertManagerInvariants(self, manager) + rssSamples.append(resourceSnapshot()['rss']) + + while manager.retiredEntryCount: + with manager._lock: + manager._cleanupRetiredLocked() + _assertManagerInvariants(self, manager) + self.assertLessEqual(maximumRetired, manager.maximumEntries) + rssValues = [value for value in rssSamples if value is not None] + rssGrowth = rssValues[-1] - rssValues[0] if len(rssValues) > 1 else None + self.report( + 'soak', + seed=seed, + operations=250_000, + duration_seconds=time.perf_counter() - started, + maximum_live=maximumLive, + maximum_retired=maximumRetired, + maximum_retired_characters=maximumRetiredCharacters, + maximum_physical_entries=maximumPhysical, + maximum_retired_batches=maximumBatches, + final_retired=manager.retiredEntryCount, + rss_samples=rssSamples, + rss_growth=rssGrowth, + ) + + def testHundredThousandEntryMergeAndClearAll(self): + """Validate a large interleaved merge and constant-work clear-all.""" + count = 100_000 + manager = LogManager( + maximumEntries=count, + maximumCharacters=count * 8, + maximumEntryCharacters=8, + autoClearEnabled=False, + ) + categories = ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + ) + for index in range(count): + manager.append(str(index), categories[index % 3]) + started = time.perf_counter_ns() + entries = manager.entries() + snapshotUs = (time.perf_counter_ns() - started) / 1_000 + self.assertEqual( + tuple(entry.sequence for entry in entries), tuple(range(1, count + 1)) + ) + watched = [] + for generation in manager._activeGenerationsLocked(): + index = _ObservedEntries(generation.entries) + generation.entries = index + watched.append(index) + started = time.perf_counter_ns() + manager.clear() + clearUs = (time.perf_counter_ns() - started) / 1_000 + self.assertEqual(sum(index.iterations for index in watched), 0) + self.assertEqual(manager.entries(), tuple()) + self.assertEqual(manager.retiredEntryCount, count) + self.report( + 'large-merge-clear', + entries=count, + snapshot_us=snapshotUs, + clear_all_us=clearUs, + synchronous_traversals=0, + ) + + @unittest.skipUnless( + os.environ.get('FURIOUS_LOG_MANAGER_BASELINE'), + 'set FURIOUS_LOG_MANAGER_BASELINE to a Git revision for comparison', + ) + def testIdenticalWorkloadAgainstPreviousImplementation(self): + """Compare the same steady, rollover, snapshot, and retention workloads.""" + revision = os.environ['FURIOUS_LOG_MANAGER_BASELINE'] + source = subprocess.run( + ['git', 'show', f'{revision}:Furious/Service/LogManager.py'], + cwd=os.getcwd(), + check=True, + capture_output=True, + text=True, + ).stdout + baselineModule = types.ModuleType('tests._baseline_log_manager') + exec( + compile(source, f'', 'exec'), baselineModule.__dict__ + ) + + def benchmark(managerClass): + results = {} + manager = managerClass(maximumEntries=50_000, autoClearEnabled=False) + started = time.perf_counter_ns() + for index in range(50_000): + manager.append(str(index), APPLICATION_LOG_CATEGORY) + results['steady_append_ms'] = (time.perf_counter_ns() - started) / 1e6 + + manager = managerClass( + maximumEntries=50_000, + autoClearMaximumEntries=20_000, + ) + for index in range(20_000): + manager.append(str(index), CORE_LOG_CATEGORY) + for index in range(20_000): + manager.append(str(index), TUN2SOCKS_LOG_CATEGORY) + started = time.perf_counter_ns() + manager.append('trigger', CORE_LOG_CATEGORY) + results['core_rollover_us'] = (time.perf_counter_ns() - started) / 1e3 + + manager = managerClass(maximumEntries=50_000, autoClearEnabled=False) + for index in range(30_000): + manager.append( + str(index), + CORE_LOG_CATEGORY if index % 5 == 0 else APPLICATION_LOG_CATEGORY, + ) + started = time.perf_counter_ns() + manager.snapshot(CORE_LOG_CATEGORY) + results['category_snapshot_us'] = (time.perf_counter_ns() - started) / 1e3 + started = time.perf_counter_ns() + manager.snapshot() + results['global_snapshot_us'] = (time.perf_counter_ns() - started) / 1e3 + + manager = managerClass(maximumEntries=1_000, autoClearEnabled=False) + started = time.perf_counter_ns() + for index in range(50_000): + manager.append(str(index), APPLICATION_LOG_CATEGORY) + results['retention_heavy_ms'] = (time.perf_counter_ns() - started) / 1e6 + return results + + current = benchmark(LogManager) + baseline = benchmark(baselineModule.LogManager) + self.assertEqual(LogManager(maximumEntries=1).entries(), tuple()) + self.report( + 'baseline-comparison', + baseline_revision=revision, + current=current, + baseline=baseline, + ratios={key: current[key] / max(baseline[key], 0.001) for key in current}, + ) + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_models_and_services.py b/tests/test_models_and_services.py index 4ee29ff7..80a4acba 100644 --- a/tests/test_models_and_services.py +++ b/tests/test_models_and_services.py @@ -46,10 +46,61 @@ from Furious.Service.MetricsHistory import ( from PySide6 import QtCore +from collections import OrderedDict import threading import unittest +import weakref -from tests.support import isolatedSettings +from tests.support import application, isolatedSettings, processQtEvents + + +class _ObservedEntryIndex(OrderedDict): + """Record global-index traversal and removal requested by one operation.""" + + def __init__(self, entries=()): + """Copy existing entries before beginning operation instrumentation.""" + self.iterationRequests = 0 + self.deletedKeys = [] + self.oldestRemovals = 0 + + super().__init__(entries) + + def __iter__(self): + """Count direct traversal of the global chronological index.""" + self.iterationRequests += 1 + + return super().__iter__() + + def items(self): + """Count item traversal of the global chronological index.""" + self.iterationRequests += 1 + + return super().items() + + def keys(self): + """Count key traversal of the global chronological index.""" + self.iterationRequests += 1 + + return super().keys() + + def values(self): + """Count value traversal of the global chronological index.""" + self.iterationRequests += 1 + + return super().values() + + def __delitem__(self, key): + """Record direct sequence-key deletion during category clearing.""" + self.deletedKeys.append(key) + + return super().__delitem__(key) + + def popitem(self, last=True): + """Record global oldest-entry eviction without counting it as a scan.""" + if not last: + self.oldestRemovals += 1 + + return super().popitem(last=last) class ProfileModelTest(unittest.TestCase): @@ -329,6 +380,66 @@ class SettingsMigrationTest(unittest.TestCase): class LogManagerTest(unittest.TestCase): """Verify bounded, categorized, and thread-safe structured logging.""" + def _assertIndexesConsistent(self, manager): + """Verify active/retired indexes and aggregate accounting agree.""" + with manager._lock: + liveSequences = [] + liveCharacters = 0 + + for generation in manager._activeGenerationsLocked(): + generationSequences = tuple(generation.entries) + categorySequences = [] + categoryCharacters = 0 + + self.assertEqual( + generationSequences, + tuple(sorted(generationSequences)), + ) + + for categoryId, categoryEntries in generation.entriesByCategory.items(): + self.assertEqual( + categoryEntries.characterCount, + sum( + len(entry.message) + for entry in categoryEntries.entries.values() + ), + ) + + for sequence, entry in categoryEntries.entries.items(): + self.assertEqual(entry.categoryId, categoryId) + self.assertIs(generation.entries[sequence], entry) + categorySequences.append(sequence) + + categoryCharacters += categoryEntries.characterCount + + self.assertEqual( + tuple(sorted(categorySequences)), + generationSequences, + ) + self.assertEqual(generation.characterCount, categoryCharacters) + liveSequences.extend(generationSequences) + liveCharacters += generation.characterCount + + self.assertEqual( + tuple(entry.sequence for entry in manager.entries()), + tuple(sorted(liveSequences)), + ) + self.assertEqual(manager.entryCount(), len(liveSequences)) + self.assertEqual(manager._retainedEntryCount, len(liveSequences)) + self.assertEqual( + manager.retainedCharacters, + liveCharacters, + ) + self.assertEqual(manager._retainedCharacters, liveCharacters) + self.assertEqual( + manager._retiredEntryCount, + sum(batch.entryCount for batch in manager._retiredBatches), + ) + self.assertEqual( + manager._retiredCharacters, + sum(batch.characterCount for batch in manager._retiredBatches), + ) + def testBoundedBufferCategoriesAndRuntimeClear(self): """Retain only the newest entries and clear runtime categories alone.""" manager = LogManager(maximumEntries=4) @@ -478,6 +589,630 @@ class LogManagerTest(unittest.TestCase): manager.entryCount(), ) + self._assertIndexesConsistent(manager) + + def testCategoryOperationsDoNotTraverseUnrelatedApplicationHistory(self): + """Make Core rollover constant regardless of runtime/history skew.""" + cases = ( + (1_000, 7_000, 2_000), + (9_000, 900, 100), + ) + + for applicationCount, coreCount, tunCount in cases: + with self.subTest( + application=applicationCount, + core=coreCount, + tun2socks=tunCount, + ): + manager = LogManager( + maximumEntries=20_000, + autoClearMaximumEntries=coreCount, + ) + + for index in range(applicationCount): + manager.append( + f'application {index}', + APPLICATION_LOG_CATEGORY, + ) + + for index in range(coreCount): + manager.append(f'core {index}', CORE_LOG_CATEGORY) + + for index in range(tunCount): + manager.append(f'tun2socks {index}', TUN2SOCKS_LOG_CATEGORY) + + applicationGeneration = manager._applicationGeneration + runtimeGeneration = manager._runtimeGeneration + observedIndexes = [] + + for generation in manager._activeGenerationsLocked(): + observedEntries = _ObservedEntryIndex(generation.entries) + generation.entries = observedEntries + observedIndexes.append(observedEntries) + + for categoryEntries in generation.entriesByCategory.values(): + observedCategory = _ObservedEntryIndex(categoryEntries.entries) + categoryEntries.entries = observedCategory + observedIndexes.append(observedCategory) + + sequence, coreEntries = manager.snapshot(CORE_LOG_CATEGORY) + + self.assertEqual( + sequence, + applicationCount + coreCount + tunCount, + ) + self.assertEqual(len(coreEntries), coreCount) + self.assertEqual( + sum(index.iterationRequests for index in observedIndexes), + 1, + ) + + for index in observedIndexes: + index.iterationRequests = 0 + + manager.append('new core after clear', CORE_LOG_CATEGORY) + + self.assertIs(manager._applicationGeneration, applicationGeneration) + self.assertIsNot(manager._runtimeGeneration, runtimeGeneration) + self.assertEqual( + sum(index.iterationRequests for index in observedIndexes), + 0, + ) + self.assertEqual( + sum(len(index.deletedKeys) for index in observedIndexes), + 0, + ) + self.assertEqual( + sum(index.oldestRemovals for index in observedIndexes), + 0, + ) + self.assertEqual( + manager.retiredEntryCount, + coreCount + tunCount, + ) + self.assertEqual( + manager.entryCount(APPLICATION_LOG_CATEGORY), + applicationCount, + ) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 1) + self.assertEqual(manager.entryCount(TUN2SOCKS_LOG_CATEGORY), 0) + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + tuple(f'application {index}' for index in range(applicationCount)) + + ('new core after clear',), + ) + self._assertIndexesConsistent(manager) + + def testWholeGenerationClearDoesNoPhysicalEntryWorkOnCaller(self): + """Swap every active stream without traversing or destroying its entries.""" + manager = LogManager(maximumEntries=1_000, autoClearEnabled=False) + otherCategory = manager.registerComponent( + 'component.persistent', + 'Persistent', + runtime=False, + ) + + for index in range(100): + manager.append(f'application {index}', APPLICATION_LOG_CATEGORY) + manager.append(f'core {index}', CORE_LOG_CATEGORY) + manager.append(f'other {index}', otherCategory.id) + + observedIndexes = [] + + for generation in manager._activeGenerationsLocked(): + observedEntries = _ObservedEntryIndex(generation.entries) + generation.entries = observedEntries + observedIndexes.append(observedEntries) + + for categoryEntries in generation.entriesByCategory.values(): + observedCategory = _ObservedEntryIndex(categoryEntries.entries) + categoryEntries.entries = observedCategory + observedIndexes.append(observedCategory) + + self.assertEqual(manager.entryCount(), 300) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 100) + self.assertGreater(manager.retainedCharacters, 0) + self.assertEqual( + sum(index.iterationRequests for index in observedIndexes), + 0, + ) + + manager.append('application after observation', APPLICATION_LOG_CATEGORY) + + self.assertEqual( + sum(index.iterationRequests for index in observedIndexes), + 0, + ) + + manager.clear() + + self.assertEqual(manager.entryCount(), 0) + self.assertEqual(manager.retainedCharacters, 0) + self.assertEqual(manager.retiredEntryCount, 301) + self.assertEqual(manager.entries(), tuple()) + self.assertEqual( + sum(index.iterationRequests for index in observedIndexes), + 0, + ) + self.assertEqual( + sum(len(index.deletedKeys) for index in observedIndexes), + 0, + ) + self.assertEqual( + sum(index.oldestRemovals for index in observedIndexes), + 0, + ) + self._assertIndexesConsistent(manager) + + def testSnapshotsNeverTraverseRetiredGenerations(self): + """Materialize only live streams after an immediate runtime rollover.""" + manager = LogManager(maximumEntries=100, autoClearEnabled=False) + manager.append('application', APPLICATION_LOG_CATEGORY) + manager.append('core', CORE_LOG_CATEGORY) + manager.append('tun2socks', TUN2SOCKS_LOG_CATEGORY) + retiredGeneration = manager._runtimeGeneration + observedGlobal = _ObservedEntryIndex(retiredGeneration.entries) + observedCore = _ObservedEntryIndex( + retiredGeneration.entriesByCategory[CORE_LOG_CATEGORY].entries + ) + retiredGeneration.entries = observedGlobal + retiredGeneration.entriesByCategory[CORE_LOG_CATEGORY].entries = observedCore + + manager.clear(runtimeOnly=True) + + self.assertEqual(manager.entries(CORE_LOG_CATEGORY), tuple()) + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('application',), + ) + self.assertEqual(observedGlobal.iterationRequests, 0) + self.assertEqual(observedCore.iterationRequests, 0) + self.assertEqual(observedGlobal.oldestRemovals, 0) + self.assertEqual(manager.retiredEntryCount, 2) + + def testSelectiveCategoryClearUsesDocumentedSharedStreamFallback(self): + """Unlink only k selected entries when another category shares a stream.""" + manager = LogManager(maximumEntries=1_000, autoClearEnabled=False) + + for index in range(100): + manager.append(f'core {index}', CORE_LOG_CATEGORY) + + for index in range(10): + manager.append(f'tun2socks {index}', TUN2SOCKS_LOG_CATEGORY) + + runtimeGeneration = manager._runtimeGeneration + observedGlobal = _ObservedEntryIndex(runtimeGeneration.entries) + observedCore = _ObservedEntryIndex( + runtimeGeneration.entriesByCategory[CORE_LOG_CATEGORY].entries + ) + runtimeGeneration.entries = observedGlobal + runtimeGeneration.entriesByCategory[CORE_LOG_CATEGORY].entries = observedCore + + manager.clear(CORE_LOG_CATEGORY) + + self.assertIs(manager._runtimeGeneration, runtimeGeneration) + self.assertEqual(observedCore.iterationRequests, 1) + self.assertEqual(len(observedGlobal.deletedKeys), 100) + self.assertEqual(observedGlobal.oldestRemovals, 0) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 0) + self.assertEqual(manager.entryCount(TUN2SOCKS_LOG_CATEGORY), 10) + self.assertEqual(manager.retiredEntryCount, 100) + + observedGlobal.iterationRequests = 0 + observedGlobal.deletedKeys.clear() + + manager.clear(TUN2SOCKS_LOG_CATEGORY) + + self.assertIsNot(manager._runtimeGeneration, runtimeGeneration) + self.assertEqual(observedGlobal.iterationRequests, 0) + self.assertEqual(observedGlobal.deletedKeys, []) + self.assertEqual(observedGlobal.oldestRemovals, 0) + self.assertEqual(manager.entryCount(), 0) + self.assertEqual(manager.retiredEntryCount, 110) + self._assertIndexesConsistent(manager) + + def testRetiredCleanupIsBoundedAndEventuallyReleasesEntries(self): + """Release at most the fixed budget from a retired generation per turn.""" + manager = LogManager(maximumEntries=100, autoClearEnabled=False) + manager.RetiredCleanupBudget = 3 + entryReferences = [] + + for index in range(10): + entry = manager.append(f'core {index}', CORE_LOG_CATEGORY) + entryReferences.append(weakref.ref(entry)) + + del entry + + runtimeGeneration = manager._runtimeGeneration + observedEntries = _ObservedEntryIndex(runtimeGeneration.entries) + runtimeGeneration.entries = observedEntries + + manager.clear(runtimeOnly=True) + + self.assertEqual(manager.retiredEntryCount, 10) + self.assertEqual(observedEntries.oldestRemovals, 0) + self.assertTrue(all(reference() is not None for reference in entryReferences)) + + manager.append('application', APPLICATION_LOG_CATEGORY) + + self.assertEqual(manager.retiredEntryCount, 7) + self.assertEqual(observedEntries.oldestRemovals, 3) + + cleanupTurns = 0 + + while manager.retiredEntryCount: + before = manager.retiredEntryCount + + with manager._lock: + cleaned = manager._cleanupRetiredLocked() + + cleanupTurns += 1 + self.assertLessEqual(cleaned, manager.RetiredCleanupBudget) + self.assertEqual(manager.retiredEntryCount, before - cleaned) + + self.assertEqual(cleanupTurns, 3) + self.assertEqual(observedEntries.oldestRemovals, 10) + self.assertEqual(manager.retiredCharacters, 0) + self.assertEqual(len(manager._retiredBatches), 0) + self.assertTrue(all(reference() is None for reference in entryReferences)) + self._assertIndexesConsistent(manager) + + def testQueuedCleanupEventuallyDrainsWithoutFurtherLogging(self): + """Finish physical reclamation through bounded queued manager turns.""" + application() + manager = LogManager(maximumEntries=100, autoClearEnabled=False) + manager.RetiredCleanupBudget = 2 + entryReferences = [] + + for index in range(9): + entry = manager.append(f'core {index}', CORE_LOG_CATEGORY) + entryReferences.append(weakref.ref(entry)) + + del entry + + manager.clear(runtimeOnly=True) + + self.assertEqual(manager.retiredEntryCount, 9) + + processQtEvents() + + self.assertEqual(manager.retiredEntryCount, 0) + self.assertEqual(manager.retiredCharacters, 0) + self.assertTrue(all(reference() is None for reference in entryReferences)) + + def testQueuedCleanupDoesNotOutliveDestroyedManager(self): + """Discard pending self-delivery at the manager's QObject boundary.""" + application() + manager = LogManager(maximumEntries=100, autoClearEnabled=False) + manager.RetiredCleanupBudget = 1 + destroyed = [] + manager.destroyed.connect(lambda: destroyed.append(True)) + + for index in range(20): + manager.append(f'core {index}', CORE_LOG_CATEGORY) + + manager.clear(runtimeOnly=True) + + managerReference = weakref.ref(manager) + manager.deleteLater() + del manager + + processQtEvents() + + self.assertEqual(destroyed, [True]) + self.assertIsNone(managerReference()) + + def testRepeatedFastClearsKeepRetiredEntryBacklogBounded(self): + """Prevent retired generations from accumulating beyond live capacity.""" + manager = LogManager(maximumEntries=128, autoClearEnabled=False) + manager.RetiredCleanupBudget = 1 + + for index in range(manager.maximumEntries): + manager.append(f'initial {index}', CORE_LOG_CATEGORY) + + manager.clear(runtimeOnly=True) + + self.assertEqual(manager.retiredEntryCount, manager.maximumEntries) + + for index in range(500): + manager.append(f'replacement {index}', CORE_LOG_CATEGORY) + manager.clear(runtimeOnly=True) + + self.assertLessEqual( + manager.retiredEntryCount, + manager.maximumEntries, + ) + self.assertLessEqual( + manager.entryCount() + manager.retiredEntryCount, + manager.maximumEntries, + ) + self.assertLessEqual( + len(manager._retiredBatches), + manager.retiredEntryCount, + ) + + while manager.retiredEntryCount: + with manager._lock: + manager._cleanupRetiredLocked() + + self.assertEqual(manager.retiredCharacters, 0) + self.assertEqual(len(manager._retiredBatches), 0) + self._assertIndexesConsistent(manager) + + def testRepeatedCoreAutoClearKeepsOnlyTheNewestRuntimeEpoch(self): + """Roll over repeatedly without exposing or accumulating older Core lines.""" + manager = LogManager( + maximumEntries=100, + autoClearMaximumEntries=1, + ) + manager.RetiredCleanupBudget = 1 + manager.append('application', APPLICATION_LOG_CATEGORY) + firstGenerationId = manager._runtimeGeneration.identifier + + for index in range(100): + manager.append(f'core {index}', CORE_LOG_CATEGORY) + + self.assertEqual( + tuple(entry.message for entry in manager.entries(CORE_LOG_CATEGORY)), + (f'core {index}',), + ) + self.assertLessEqual(manager.retiredEntryCount, 1) + + self.assertGreater(manager._runtimeGeneration.identifier, firstGenerationId) + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('application', 'core 99'), + ) + self.assertEqual(manager.entryCount(APPLICATION_LOG_CATEGORY), 1) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 1) + self._assertIndexesConsistent(manager) + + def testRetentionEvictsOldestEntriesFromBothIndexes(self): + """Keep category indexes synchronized through repeated global eviction.""" + manager = LogManager( + maximumEntries=4, + maximumCharacters=20, + maximumEntryCharacters=20, + autoClearEnabled=False, + ) + categoryIds = ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + TUN2SOCKS_LOG_CATEGORY, + ) + + for index in range(3): + manager.append(f'{index:04d}', categoryIds[index]) + + observedIndexes = [] + + for generation in manager._activeGenerationsLocked(): + observedEntries = _ObservedEntryIndex(generation.entries) + generation.entries = observedEntries + observedIndexes.append(observedEntries) + + for index in range(3, 20): + manager.append(f'{index:04d}', categoryIds[index % len(categoryIds)]) + + entries = manager.entries() + + self.assertEqual( + tuple(entry.message for entry in entries), + ('0016', '0017', '0018', '0019'), + ) + self.assertEqual( + sum(index.oldestRemovals for index in observedIndexes), + 16, + ) + self.assertEqual(manager.retainedCharacters, 16) + self._assertIndexesConsistent(manager) + + def testGlobalRetentionIgnoresRetiredGenerations(self): + """Evict the oldest live stream head across a runtime rollover.""" + manager = LogManager( + maximumEntries=4, + maximumCharacters=100, + maximumEntryCharacters=100, + autoClearEnabled=False, + ) + manager.RetiredCleanupBudget = 1 + + manager.append('application 1', APPLICATION_LOG_CATEGORY) + manager.append('old core', CORE_LOG_CATEGORY) + manager.append('application 2', APPLICATION_LOG_CATEGORY) + manager.append('old tun2socks', TUN2SOCKS_LOG_CATEGORY) + + manager.clear(runtimeOnly=True) + + self.assertEqual(manager.retiredEntryCount, 2) + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('application 1', 'application 2'), + ) + + manager.append('new core', CORE_LOG_CATEGORY) + manager.append('application 3', APPLICATION_LOG_CATEGORY) + manager.append('new tun2socks', TUN2SOCKS_LOG_CATEGORY) + + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ( + 'application 2', + 'new core', + 'application 3', + 'new tun2socks', + ), + ) + self.assertEqual(manager.retiredEntryCount, 0) + self.assertEqual(manager.entryCount(APPLICATION_LOG_CATEGORY), 2) + self.assertEqual(manager.entryCount(CORE_LOG_CATEGORY), 1) + self.assertEqual(manager.entryCount(TUN2SOCKS_LOG_CATEGORY), 1) + self._assertIndexesConsistent(manager) + + def testRegisteredCategoryClearAndClearAllPreserveSignals(self): + """Clear one late category and then all entries with compatible signals.""" + manager = LogManager(maximumEntries=20, autoClearEnabled=False) + category = manager.registerComponent( + 'component.audit', + 'Audit', + runtime=False, + ) + cleared = [] + manager.entriesCleared.connect(cleared.append) + + manager.append('application', APPLICATION_LOG_CATEGORY) + manager.append('audit 1', category.id) + manager.append('core', CORE_LOG_CATEGORY) + manager.append('audit 2', category.id) + + self.assertEqual(manager.snapshot('unknown.category')[1], tuple()) + self.assertEqual(manager.entryCount(category.id), 2) + self.assertEqual(manager.plainText(category.id), 'audit 1\naudit 2') + + manager.clear(runtimeOnly=True) + + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('application', 'audit 1', 'audit 2'), + ) + self.assertEqual( + cleared, + [frozenset({CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY})], + ) + + manager.clear(category.id) + + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('application',), + ) + self.assertEqual( + cleared, + [ + frozenset({CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY}), + frozenset({category.id}), + ], + ) + self._assertIndexesConsistent(manager) + + manager.clear() + manager.clear() + + self.assertEqual(manager.entries(), tuple()) + self.assertEqual( + cleared, + [ + frozenset({CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY}), + frozenset({category.id}), + None, + ], + ) + self._assertIndexesConsistent(manager) + + def testConcurrentAppendClearAndSnapshotsKeepIndexesConsistent(self): + """Serialize snapshots and category clears racing with log producers.""" + manager = LogManager(maximumEntries=5_000, autoClearEnabled=False) + category = manager.registerComponent('component.concurrent', 'Concurrent') + started = threading.Event() + finished = threading.Event() + failures = [] + observations = [] + + def produce(): + """Append interleaved categories while the observer reads and clears.""" + try: + started.set() + + categoryIds = ( + APPLICATION_LOG_CATEGORY, + CORE_LOG_CATEGORY, + category.id, + ) + + for index in range(3_000): + manager.append( + f'entry {index}', + categoryIds[index % len(categoryIds)], + ) + except Exception as error: + failures.append(error) + finally: + finished.set() + + def observeAndClear(): + """Take filtered snapshots and clear one category during ingestion.""" + try: + started.wait(5) + + while True: + manager.snapshot(CORE_LOG_CATEGORY) + manager.snapshot(category.id) + manager.clear(category.id) + observations.append(True) + + if finished.wait(0.0001): + break + except Exception as error: + failures.append(error) + + observer = threading.Thread(target=observeAndClear) + producer = threading.Thread(target=produce) + observer.start() + producer.start() + producer.join(10) + observer.join(10) + + self.assertFalse(producer.is_alive()) + self.assertFalse(observer.is_alive()) + self.assertEqual(failures, []) + self.assertTrue(observations) + + sequences = tuple(entry.sequence for entry in manager.entries()) + + self.assertEqual(sequences, tuple(sorted(sequences))) + self.assertEqual(len(sequences), len(set(sequences))) + self._assertIndexesConsistent(manager) + + def testEntrySignalsAndCrossThreadChangeNotificationsRemainCompatible(self): + """Emit every entry and coalesce presentation refreshes by producer batch.""" + application() + manager = LogManager(maximumEntries=100, autoClearEnabled=False) + added = [] + changed = [] + manager.entryAdded.connect(added.append) + manager.entriesChanged.connect(changed.append) + + def produce(): + """Publish one batch without allowing the Qt queue to drain midway.""" + for index in range(50): + manager.append(f'worker {index}', CORE_LOG_CATEGORY) + + worker = threading.Thread(target=produce) + worker.start() + worker.join(5) + + self.assertFalse(worker.is_alive()) + self.assertEqual(added, []) + self.assertEqual(changed, []) + + processQtEvents() + + self.assertEqual(len(added), 50) + self.assertEqual( + tuple(entry.message for entry in added), + tuple(f'worker {index}' for index in range(50)), + ) + self.assertEqual(changed, [50]) + + manager.append('application 1', APPLICATION_LOG_CATEGORY) + manager.append('application 2', APPLICATION_LOG_CATEGORY) + + self.assertEqual(len(added), 52) + + processQtEvents() + + self.assertEqual(changed, [50, 52]) + class MetricsHistoryTest(unittest.TestCase): """Verify bounded history and metric-specific aggregation semantics."""