Implement generation-based log storage

This commit is contained in:
Loren Eteval
2026-08-27 14:58:38 +08:00
parent 89b18fd01e
commit fbd28c48d4
5 changed files with 2596 additions and 84 deletions
+420 -83
View File
@@ -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(
+6
View File
@@ -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)
+4
View File
@@ -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
File diff suppressed because it is too large Load Diff
+736 -1
View File
@@ -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."""