diff --git a/Furious/Backends/ExternalCore/Process.py b/Furious/Backends/ExternalCore/Process.py index f019032..db1f04f 100644 --- a/Furious/Backends/ExternalCore/Process.py +++ b/Furious/Backends/ExternalCore/Process.py @@ -105,6 +105,28 @@ class ExternalCoreProcess(CoreRuntime): logger.exception('external core output callback failed') + def _emitOutputs(self, messages): + """Forward one decoded OS-read batch through an optional batch route.""" + messages = tuple(message for message in messages if message) + + if not messages or not callable(self._messageCallback): + return + + appendMany = getattr(self._messageCallback, 'appendMany', None) + + if not callable(appendMany): + for message in messages: + self._emitOutput(message) + + return + + try: + appendMany(messages) + except Exception: + # Any non-exit exceptions + + logger.exception('external core output callback failed') + def _readStream(self, stream: BinaryIO, label: str): """Drain one child pipe until EOF without retaining process output.""" pending = b'' @@ -119,6 +141,7 @@ class ExternalCoreProcess(CoreRuntime): pending += chunk lines = pending.split(b'\n') pending = lines.pop() + messages = [] for line in lines: message = line.decode('utf-8', 'replace').rstrip('\r') @@ -126,10 +149,12 @@ class ExternalCoreProcess(CoreRuntime): if not message: continue - self._emitOutput( + messages.append( message if label == 'stdout' else f'[stderr] {message}' ) + self._emitOutputs(messages) + if len(pending) >= self.MaximumPendingOutput: message = pending.decode('utf-8', 'replace') diff --git a/Furious/Controllers/ConnectionController.py b/Furious/Controllers/ConnectionController.py index 4c61d40..3424f96 100644 --- a/Furious/Controllers/ConnectionController.py +++ b/Furious/Controllers/ConnectionController.py @@ -26,10 +26,10 @@ from Furious.Plugins import getPluginRegistry from Furious.Qt.DynamicTranslate import gettext as _ from Furious.Repository import Storage from Furious.Service import ( + CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY, ConnectionManager, UpdateManager, - coreLogCallback, ) from PySide6 import QtCore @@ -282,7 +282,7 @@ class ConnectionController(QtCore.QObject): configuration, routing=AppSettings.get('Routing'), exitCallback=self.coreExitCallback, - msgCallbackCore=coreLogCallback(logManager), + msgCallbackCore=logManager.callback(CORE_LOG_CATEGORY), msgCallbackTUN_=logManager.callback( TUN2SOCKS_LOG_CATEGORY, source='Tun2socks', diff --git a/Furious/Core/CoreProcessWorker.py b/Furious/Core/CoreProcessWorker.py index 4c1b29f..996e1ab 100644 --- a/Furious/Core/CoreProcessWorker.py +++ b/Furious/Core/CoreProcessWorker.py @@ -212,6 +212,8 @@ class MsgQueue(multiprocessing.queues.Queue): return hasMessages = False + appendMany = getattr(self.callback, 'appendMany', None) + messages = [] if callable(appendMany) else None for _ in range(self.MAXIMUM_MESSAGES_PER_TICK): msg = self.getNoWait() @@ -222,7 +224,13 @@ class MsgQueue(multiprocessing.queues.Queue): hasMessages = True if not msg.isspace(): - self.callback(msg) + if messages is None: + self.callback(msg) + else: + messages.append(msg) + + if messages: + appendMany(messages) if hasMessages: nextTimeout = self.ACTIVE_DRAIN_INTERVAL diff --git a/Furious/Service/LogManager.py b/Furious/Service/LogManager.py index 00b61d2..6673e1d 100644 --- a/Furious/Service/LogManager.py +++ b/Furious/Service/LogManager.py @@ -24,6 +24,7 @@ from Furious.Models.Logging import LogCategory, LogEntry from PySide6 import QtCore from collections import OrderedDict, deque +from dataclasses import dataclass from datetime import datetime from heapq import merge from typing import Optional @@ -37,8 +38,9 @@ __all__ = [ 'CORE_LOG_CATEGORY', 'TUN2SOCKS_LOG_CATEGORY', 'ApplicationLogHandler', + 'LogCursor', + 'LogEntryBatch', 'LogManager', - 'coreLogCallback', 'formatLogEntry', ] @@ -48,11 +50,78 @@ CORE_LOG_CATEGORY = 'core' TUN2SOCKS_LOG_CATEGORY = 'component.tun2socks' +@dataclass(frozen=True) +class LogCursor: + """Opaque synchronization position for one filtered live-log view. + + Instances are issued by :meth:`LogManager.entriesSince`. Callers should + retain and pass them back unchanged rather than construct them or interpret + their fields. In particular, ``generationStates`` encodes private generation + identities and revisions whose representation may change independently of + this API. + """ + + sequence: int + categoryId: str + generationStates: tuple[tuple[int, int], ...] + + +@dataclass(frozen=True) +class LogEntryBatch: + """Describe one incremental suffix or a required full resynchronization.""" + + cursor: LogCursor + entries: tuple[LogEntry, ...] + firstRetainedSequence: Optional[int] + resetRequired: bool = False + + def formatLogEntry(entry: LogEntry) -> str: """Return the producer-formatted text stored by one structured entry.""" return entry.message +class _LogCallback: + """Adapt one fixed log route to single-line and batch producers.""" + + __slots__ = ('manager', 'categoryId', 'source', 'severity') + + def __init__(self, manager, categoryId: str, source: str, severity: str): + """Capture one application-lifetime manager route.""" + self.manager = manager + self.categoryId = categoryId + self.source = source + self.severity = severity + + def __call__(self, line): + """Append one line safely and return its entry, or ``None`` on failure.""" + try: + return self.manager.append( + line, + self.categoryId, + source=self.source, + severity=self.severity, + ) + except Exception: + # Any non-exit exceptions + + return None + + def appendMany(self, lines): + """Append one natural producer batch under the same safe contract.""" + try: + return self.manager.appendMany( + lines, + self.categoryId, + source=self.source, + severity=self.severity, + ) + except Exception: + # Any non-exit exceptions + + return tuple() + + class _CategoryEntries: """Own one category's ordered entries and aggregate character count.""" @@ -98,6 +167,7 @@ class _EntryGeneration: 'entries', 'entriesByCategory', 'characterCount', + 'revision', ) def __init__(self, identifier: int, scope: str): @@ -107,6 +177,7 @@ class _EntryGeneration: self.entries: OrderedDict[int, LogEntry] = OrderedDict() self.entriesByCategory: dict[str, _CategoryEntries] = {} self.characterCount = 0 + self.revision = 0 @property def entryCount(self) -> int: @@ -123,16 +194,20 @@ class _EntryGeneration: 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: @@ -161,6 +236,7 @@ class _EntryGeneration: del self.entries[sequence] self.characterCount -= categoryEntries.characterCount + self.revision += 1 return categoryEntries @@ -184,11 +260,12 @@ class LogManager(QtCore.QObject): 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. + # Append-side cleanup charges one small unit per accepted input to keep + # producer latency stable. Queued Qt cleanup uses a larger bounded turn to + # drain idle backlogs. Together they bound live-plus-retired entries under + # repeated fast clears without making a rollover release its own generation. RetiredCleanupBudget = 64 + AppendRetiredCleanupBudget = 1 TruncationMarker = '\n... [log entry truncated]' categoryRegistered = QtCore.Signal(object) @@ -404,11 +481,13 @@ class LogManager(QtCore.QObject): 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) @@ -423,7 +502,11 @@ class LogManager(QtCore.QObject): return bool(removedEntries) def _cleanupRetiredLocked(self, budget: Optional[int] = None) -> int: - """Physically release at most one fixed batch of logically dead entries.""" + """Release a bounded number of retired entries in global FIFO order. + + The bound counts entries, not retired batches, and may therefore cross + batch boundaries. ``RetiredCleanupBudget`` supplies the default bound. + """ remaining = ( max(1, int(self.RetiredCleanupBudget)) if budget is None @@ -434,8 +517,10 @@ class LogManager(QtCore.QObject): while remaining and self._retiredBatches: batch = self._retiredBatches[0] characterCount = batch.discardOldestPhysical() + self._retiredEntryCount -= 1 self._retiredCharacters -= characterCount + cleaned += 1 remaining -= 1 @@ -454,6 +539,7 @@ class LogManager(QtCore.QObject): with self._lock: if self._retiredBatches and not self._retiredCleanupPending: self._retiredCleanupPending = True + shouldNotify = True if shouldNotify: @@ -465,6 +551,7 @@ class LogManager(QtCore.QObject): with self._lock: self._retiredCleanupPending = False self._cleanupRetiredLocked() + shouldContinue = bool(self._retiredBatches) if shouldContinue: @@ -486,6 +573,7 @@ class LogManager(QtCore.QObject): key=lambda generation: next(iter(generation.entries)), ) entry = oldestGeneration.removeOldest() + self._retainedEntryCount -= 1 self._retainedCharacters -= len(entry.message) @@ -511,6 +599,128 @@ class LogManager(QtCore.QObject): return text[: self._maximumEntryCharacters - len(marker)] + marker + @staticmethod + def _normalizedCategoryId(categoryId: Optional[str]) -> str: + """Return the canonical filter identifier used by incremental cursors.""" + return ALL_LOGS_FILTER if categoryId in (None, ALL_LOGS_FILTER) else categoryId + + def _generationStatesLocked( + self, + categoryId: str, + ) -> tuple[tuple[int, int], ...]: + """Return the active structural state relevant to one filtered view.""" + if categoryId == ALL_LOGS_FILTER: + return tuple( + (generation.identifier, generation.revision) + for generation in self._activeGenerationsLocked() + ) + + category = self._categories.get(categoryId) + + if category is None: + return tuple() + + generation = self._generationForCategoryLocked(category) + + return ((generation.identifier, generation.revision),) + + def _entriesLocked(self, categoryId: str) -> tuple[LogEntry, ...]: + """Return the complete active entries for one canonical filter.""" + if categoryId == ALL_LOGS_FILTER: + return tuple( + merge( + *( + generation.entries.values() + for generation in self._activeGenerationsLocked() + ), + key=lambda entry: entry.sequence, + ) + ) + + category = self._categories.get(categoryId) + + if category is None: + return tuple() + + generation = self._generationForCategoryLocked(category) + categoryEntries = generation.categoryEntries(categoryId) + + return ( + tuple(categoryEntries.entries.values()) + if categoryEntries is not None + else tuple() + ) + + @staticmethod + def _orderedEntriesAfter(entries, sequence: int) -> tuple[LogEntry, ...]: + """Read only the suffix newer than *sequence* from one ordered index.""" + suffix = [] + + for entry in reversed(entries.values()): + if entry.sequence <= sequence: + break + + suffix.append(entry) + + suffix.reverse() + + return tuple(suffix) + + def _entriesAfterLocked( + self, + sequence: int, + categoryId: str, + ) -> tuple[LogEntry, ...]: + """Return only active entries newer than one valid cursor sequence.""" + if categoryId == ALL_LOGS_FILTER: + return tuple( + merge( + *( + self._orderedEntriesAfter(generation.entries, sequence) + for generation in self._activeGenerationsLocked() + ), + key=lambda entry: entry.sequence, + ) + ) + + category = self._categories.get(categoryId) + + if category is None: + return tuple() + + generation = self._generationForCategoryLocked(category) + categoryEntries = generation.categoryEntries(categoryId) + + return ( + self._orderedEntriesAfter(categoryEntries.entries, sequence) + if categoryEntries is not None + else tuple() + ) + + def _firstRetainedSequenceLocked(self, categoryId: str) -> Optional[int]: + """Return the oldest sequence still visible through one filter.""" + if categoryId == ALL_LOGS_FILTER: + heads = ( + next(iter(generation.entries)) + for generation in self._activeGenerationsLocked() + if generation.entryCount + ) + + return min(heads, default=None) + + category = self._categories.get(categoryId) + + if category is None: + return None + + generation = self._generationForCategoryLocked(category) + categoryEntries = generation.categoryEntries(categoryId) + + if categoryEntries is None: + return None + + return next(iter(categoryEntries.entries)) + def _markEntriesChangedLocked(self) -> bool: """Mark one coalesced presentation notification while already locked.""" if self._changeNotificationPending: @@ -635,8 +845,39 @@ class LogManager(QtCore.QObject): source: str = '', severity: str = '', ) -> LogEntry: - """Append with fixed cleanup work and notify interested presenters.""" - clearedCategoryIds = frozenset() + """Append one entry while preserving the compatibility signal contract.""" + return self.appendMany( + (message,), + categoryId, + timestamp=timestamp, + source=source, + severity=severity, + )[0] + + def appendMany( + self, + messages, + categoryId: str = APPLICATION_LOG_CATEGORY, + *, + timestamp: Optional[datetime] = None, + source: str = '', + severity: str = '', + ) -> tuple[LogEntry, ...]: + """Atomically append a batch, then emit compatibility events in order. + + ``entryAdded`` observers see the fully committed batch, never an + intermediate collection state. Presenters should use ``entriesChanged`` + and ``entriesSince`` instead of treating ``entryAdded`` as UI state. + """ + messages = tuple(messages) + + if not messages: + return tuple() + + if timestamp is not None and not isinstance(timestamp, datetime): + raise TypeError('log timestamp must be a datetime') + + events = [] shouldNotifyEntriesChanged = False retiredCleanupNeeded = False @@ -646,56 +887,70 @@ 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): - raise TypeError('log timestamp must be a datetime') - - self._sequence += 1 - - entry = LogEntry( - message=self._normalizeMessage(message), - timestamp=timestamp, - categoryId=category.id, - categoryLabel=category.displayName, - categoryTranslatable=category.translatable, - source=str(source) if source else '', - severity=str(severity) if severity else '', - sequence=self._sequence, + # Validate every caller-controlled conversion before changing the + # collection so a malformed item cannot leave a silent partial batch. + normalizedMessages = tuple( + self._normalizeMessage(message) for message in messages ) + normalizedSource = str(source) if source else '' + normalizedSeverity = str(severity) if severity else '' - if ( - category.id == CORE_LOG_CATEGORY - and self._autoClearEnabled - 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 two replaceable streams that roll over without touching - # their entries or performing synchronous object destruction. - clearedCategoryIds = self._nonApplicationCategoryIds - self._clearNonApplicationLocked() + for normalizedMessage in normalizedMessages: + # Reclaim before accepting the next input so a Core-triggering + # append itself retains O(1) generation-rollover latency. A + # tiny per-input charge still bounds backlog across large batches. + if self._retiredBatches: + self._cleanupRetiredLocked( + max(1, int(self.AppendRetiredCleanupBudget)) + ) + + entryTimestamp = datetime.now() if timestamp is None else timestamp + self._sequence += 1 + + entry = LogEntry( + message=normalizedMessage, + timestamp=entryTimestamp, + categoryId=category.id, + categoryLabel=category.displayName, + categoryTranslatable=category.translatable, + source=normalizedSource, + severity=normalizedSeverity, + sequence=self._sequence, + ) + clearedCategoryIds = frozenset() + + if ( + category.id == CORE_LOG_CATEGORY + and self._autoClearEnabled + and self._categoryEntryCountLocked(CORE_LOG_CATEGORY) + >= self._autoClearMaximumEntries + ): + # Preserve repeated-append Core rollover semantics within + # the batch, including the clear-before-entry signal order. + clearedCategoryIds = self._nonApplicationCategoryIds + + self._clearNonApplicationLocked() + + characterCount = len(entry.message) + generation = self._generationForCategoryLocked(category) + generation.append(entry, characterCount) + + self._retainedEntryCount += 1 + self._retainedCharacters += characterCount + self._enforceRetentionLimitsLocked() + + events.append((clearedCategoryIds, entry)) - 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(clearedCategoryIds) + for clearedCategoryIds, entry in events: + if clearedCategoryIds: + self.entriesCleared.emit(clearedCategoryIds) - self.entryAdded.emit(entry) + # Kept for compatibility and non-presentation observers. LogPage + # deliberately pulls a coalesced batch through entriesSince(). + self.entryAdded.emit(entry) if shouldNotifyEntriesChanged: self._entriesChangedRequested.emit() @@ -703,7 +958,7 @@ class LogManager(QtCore.QObject): if retiredCleanupNeeded: self._requestRetiredCleanup() - return entry + return tuple(entry for _clearedCategoryIds, entry in events) def callback( self, @@ -712,23 +967,8 @@ class LogManager(QtCore.QObject): source: str = '', severity: str = '', ): - """Return a safe line callback for a process or another log producer.""" - - def appendLine(line): - """Append one externally produced line without affecting its producer.""" - try: - self.append( - line, - categoryId, - source=source, - severity=severity, - ) - except Exception: - # Any non-exit exceptions - - pass - - return appendLine + """Return a safe callable supporting single lines and natural batches.""" + return _LogCallback(self, categoryId, source, severity) def entries(self, categoryId: Optional[str] = None) -> tuple[LogEntry, ...]: """Return an immutable snapshot, optionally filtered by category.""" @@ -740,34 +980,58 @@ class LogManager(QtCore.QObject): ) -> tuple[int, tuple[LogEntry, ...]]: """Return an O(n) global or O(k) category-specific immutable snapshot.""" with self._lock: - if categoryId in (None, ALL_LOGS_FILTER): - # 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( - merge( - *( - generation.entries.values() - for generation in self._activeGenerationsLocked() - ), - key=lambda entry: entry.sequence, - ) - ) - else: - category = self._categories.get(categoryId) + return ( + self._sequence, + self._entriesLocked(self._normalizedCategoryId(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() - ) + def entriesSince( + self, + cursor: Optional[LogCursor] = None, + categoryId: Optional[str] = None, + ) -> LogEntryBatch: + """Return one atomic filtered suffix and its next synchronization cursor. - return self._sequence, entries + Treat ``LogEntryBatch.cursor`` as opaque and pass it back unchanged with + the same filter on the next call. + + A missing or structurally invalid cursor returns the complete retained + view with ``resetRequired`` set. Retention-only prefix eviction keeps the + cursor valid; ``firstRetainedSequence`` lets a presenter prune exactly + that obsolete prefix without rebuilding or copying retained history. + """ + if cursor is not None and not isinstance(cursor, LogCursor): + raise TypeError('cursor must be a LogCursor or None') + + normalizedCategoryId = self._normalizedCategoryId(categoryId) + + with self._lock: + generationStates = self._generationStatesLocked(normalizedCategoryId) + resetRequired = ( + cursor is None + or cursor.categoryId != normalizedCategoryId + or cursor.generationStates != generationStates + or cursor.sequence > self._sequence + ) + entries = ( + self._entriesLocked(normalizedCategoryId) + if resetRequired + else self._entriesAfterLocked(cursor.sequence, normalizedCategoryId) + ) + nextCursor = LogCursor( + sequence=self._sequence, + categoryId=normalizedCategoryId, + generationStates=generationStates, + ) + + return LogEntryBatch( + cursor=nextCursor, + entries=entries, + firstRetainedSequence=self._firstRetainedSequenceLocked( + normalizedCategoryId + ), + resetRequired=resetRequired, + ) def clear( self, @@ -863,8 +1127,3 @@ class ApplicationLogHandler(logging.Handler): # Any non-exit exceptions self.handleError(record) - - -def coreLogCallback(manager: LogManager): - """Create a callback for output from the currently selected proxy core.""" - return manager.callback(CORE_LOG_CATEGORY) diff --git a/Furious/Service/__init__.py b/Furious/Service/__init__.py index 38d9421..46ef01a 100644 --- a/Furious/Service/__init__.py +++ b/Furious/Service/__init__.py @@ -36,8 +36,9 @@ from .LogManager import ( CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY, ApplicationLogHandler, + LogCursor, + LogEntryBatch, LogManager, - coreLogCallback, formatLogEntry, ) from .MetricsHistory import ( @@ -86,6 +87,8 @@ __all__ = [ 'CORE_LOG_CATEGORY', 'TUN2SOCKS_LOG_CATEGORY', 'ApplicationLogHandler', + 'LogCursor', + 'LogEntryBatch', 'LogManager', 'DOWNLOAD_SPEED_METRIC', 'DOWNLOAD_USAGE_METRIC', @@ -95,7 +98,6 @@ __all__ = [ 'MetricSample', 'MetricsHistory', 'PluginNavigationManager', - 'coreLogCallback', 'SubscriptionImportResult', 'SubscriptionImportService', 'SubscriptionSource', diff --git a/Furious/Widget/ServerTableView.py b/Furious/Widget/ServerTableView.py index afe28b1..8ed8b65 100644 --- a/Furious/Widget/ServerTableView.py +++ b/Furious/Widget/ServerTableView.py @@ -33,10 +33,10 @@ from Furious.Qt import * from Furious.Qt.Signals import connectWeakly from Furious.Qt import gettext as _ from Furious.Service import ( + CORE_LOG_CATEGORY, ConnectionManager, SubscriptionManager, SubscriptionUpdateBatch, - coreLogCallback, ) from Furious.Widget.WaitingSpinner import WaitingSpinner @@ -346,7 +346,7 @@ class TestDownloadSpeedWorker(HttpGetManager): configcopy, AppBuiltinRouting.Global.value, self.coreExitCallback, - msgCallbackCore=coreLogCallback(AppLogManager()), + msgCallbackCore=AppLogManager().callback(CORE_LOG_CATEGORY), deepcopy=False, proxyModeOnly=True, log=False, diff --git a/Furious/Window/LogPage.py b/Furious/Window/LogPage.py index cd75f0b..bb6283e 100644 --- a/Furious/Window/LogPage.py +++ b/Furious/Window/LogPage.py @@ -40,10 +40,15 @@ from PySide6 import QtCore, QtGui from PySide6.QtGui import QTextCursor from PySide6.QtWidgets import * -from bisect import bisect_left +from collections import deque +from dataclasses import dataclass + +import logging __all__ = ['LogPage'] +logger = logging.getLogger(__name__) + registerAppSettings('LogViewerWidgetPointSize') registerAppSettings('LogViewerSelectedCategory', default=ALL_LOGS_FILTER) @@ -54,6 +59,44 @@ _LEGACY_POINT_SIZE_SETTINGS = ( ) +@dataclass(frozen=True) +class _RenderedEntryMetadata: + """Track one normalized entry's sequence and owned document paragraphs. + + LogManager strips terminal CR/LF before entries are joined by one newline. + The first entry's base block owns the document's initial block; each later + entry's base block owns its preceding inter-entry separator. Any additional + blocks counted in ``blockCount`` belong to paragraph breaks inside that + entry, so the per-entry counts remain additive for prefix removal. + """ + + sequence: int + blockCount: int + + +class _DocumentBlockAccountingError(RuntimeError): + """Report a recoverable mismatch between rendered metadata and Qt blocks.""" + + +def _documentBlockCount(text: str) -> int: + """Return the QTextDocument block count produced by plain *text*.""" + blockCount = 1 + text.count('\n') + text.count('\u2029') + + for index, character in enumerate(text): + if character == '\r' and (index + 1 == len(text) or text[index + 1] != '\n'): + blockCount += 1 + + return blockCount + + +def _renderedEntryMetadata(entry) -> _RenderedEntryMetadata: + """Return immutable block ownership for one normalized log entry.""" + return _RenderedEntryMetadata( + sequence=entry.sequence, + blockCount=_documentBlockCount(formatLogEntry(entry)), + ) + + def _migratePointSizeSettings(): """Preserve one legacy size and remove the obsolete per-source settings.""" settings = QtCore.QSettings() @@ -165,7 +208,6 @@ class LogPage(Mixins.QTranslatable, QMainWindow): ) self.textBrowser.setLineWrapMode(DraculaTextBrowser.LineWrapMode.NoWrap) self.textBrowser.setUndoRedoEnabled(False) - self.textBrowser.document().setMaximumBlockCount(manager.maximumEntries) self.highlightOverlay = QFrame(self.textBrowser.viewport()) self.highlightOverlay.setObjectName('LogHighlightOverlay') @@ -198,8 +240,9 @@ class LogPage(Mixins.QTranslatable, QMainWindow): self._entriesDirty = True self._representationInvalid = True self._renderedSequence = 0 + self._entryCursor = None self._renderedCategoryId = None - self._renderedEntrySequences = tuple() + self._renderedEntries = deque() self._followTail = self._autoScrollDown self._documentMutation = False self._highlightNextBlock = None @@ -347,9 +390,8 @@ class LogPage(Mixins.QTranslatable, QMainWindow): # 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. + # entriesChanged coalesces producers, then one cursor query returns only + # the missing suffix plus an eviction boundary for prefix pruning. self.manager.entriesChanged.connect(self._entriesChanged) self.manager.entriesCleared.connect(self._entriesCleared) @@ -450,31 +492,6 @@ class LogPage(Mixins.QTranslatable, QMainWindow): self.highlightSpinner.stop() self.highlightOverlay.hide() - @staticmethod - def _incrementalPlan(oldSequences, currentSequences): - """Return leading removals and the first new snapshot entry. - - Retention pruning can only remove a prefix, while ordinary collection - can only append a suffix. Any other relationship requires one clean - rebuild because a clear/filter change invalidated the representation. - """ - if not oldSequences: - return 0, 0 - - if not currentSequences: - return len(oldSequences), 0 - - oldStart = bisect_left(oldSequences, currentSequences[0]) - survivors = oldSequences[oldStart:] - - if len(currentSequences) < len(survivors): - return None - - if tuple(currentSequences[: len(survivors)]) != survivors: - return None - - return oldStart, len(survivors) - def _scrollRatio(self) -> float: """Return the current vertical position as a bounded range ratio.""" scrollbar = self.textBrowser.verticalScrollBar() @@ -484,23 +501,23 @@ class LogPage(Mixins.QTranslatable, QMainWindow): return min(1.0, max(0.0, scrollbar.value() / scrollbar.maximum())) - def _removeLeadingBlocks(self, count: int): - """Remove *count* rendered entries in one document edit.""" - if count <= 0: + def _removeLeadingDocumentBlocks(self, blockCount: int, *, removeAll=False): + """Remove an exact owned paragraph prefix in one document edit.""" + if blockCount <= 0: return - if count >= len(self._renderedEntrySequences): + if removeAll: self.textBrowser.clear() return document = self.textBrowser.document() - firstRetainedBlock = document.findBlockByNumber(count) + firstRetainedBlock = document.findBlockByNumber(blockCount) if not firstRetainedBlock.isValid(): - self.textBrowser.clear() - - return + raise _DocumentBlockAccountingError( + 'rendered log block accounting diverged' + ) cursor = QTextCursor(document) cursor.setPosition(0) @@ -547,20 +564,13 @@ class LogPage(Mixins.QTranslatable, QMainWindow): return 0 if entries else None - def _synchronizeDocument(self, entries, categoryId: str): - """Apply one snapshot with the smallest safe document mutation.""" - currentSequences = tuple(entry.sequence for entry in entries) + def _synchronizeDocument(self, batch, categoryId: str): + """Apply one incremental manager batch with the smallest safe mutation.""" fullRebuild = ( - self._representationInvalid or self._renderedCategoryId != categoryId + self._representationInvalid + or self._renderedCategoryId != categoryId + or batch.resetRequired ) - plan = None - - if not fullRebuild: - plan = self._incrementalPlan( - self._renderedEntrySequences, - currentSequences, - ) - fullRebuild = plan is None oldRatio = self._scrollRatio() firstHighlightBlock = None @@ -573,26 +583,57 @@ class LogPage(Mixins.QTranslatable, QMainWindow): try: if fullRebuild: - firstHighlightBlock = self._replaceEntries(entries) - else: - dropCount, firstNewEntry = plan + firstHighlightBlock = self._replaceEntries(batch.entries) - if self._highlightNextBlock is not None and dropCount: + self._renderedEntries = deque( + _renderedEntryMetadata(entry) for entry in batch.entries + ) + else: + firstRetainedSequence = batch.firstRetainedSequence + + if firstRetainedSequence is None: + dropEntryCount = len(self._renderedEntries) + droppedBlockCount = sum( + metadata.blockCount for metadata in self._renderedEntries + ) + else: + dropEntryCount = 0 + droppedBlockCount = 0 + + for metadata in self._renderedEntries: + if metadata.sequence >= firstRetainedSequence: + break + + dropEntryCount += 1 + droppedBlockCount += metadata.blockCount + + self._removeLeadingDocumentBlocks( + droppedBlockCount, + removeAll=bool(dropEntryCount) + and dropEntryCount == len(self._renderedEntries), + ) + + for _index in range(dropEntryCount): + self._renderedEntries.popleft() + + if self._highlightNextBlock is not None and droppedBlockCount: self._highlightNextBlock = max( 0, - self._highlightNextBlock - dropCount, + self._highlightNextBlock - droppedBlockCount, ) - self._removeLeadingBlocks(dropCount) + leadingBoundaryChanged = bool(dropEntryCount) - leadingBoundaryChanged = bool(dropCount) - - newEntries = entries[firstNewEntry:] firstHighlightBlock = self._appendEntries( - newEntries, - hasExistingEntries=bool(currentSequences[:firstNewEntry]), + batch.entries, + hasExistingEntries=bool(self._renderedEntries), ) - restoreManualPosition = bool(dropCount) + + self._renderedEntries.extend( + _renderedEntryMetadata(entry) for entry in batch.entries + ) + + restoreManualPosition = bool(dropEntryCount) finally: self.textBrowser.setSyntaxHighlightingEnabled(True) @@ -609,7 +650,6 @@ class LogPage(Mixins.QTranslatable, QMainWindow): self.textBrowser.viewport().update() self._renderedCategoryId = categoryId - self._renderedEntrySequences = currentSequences self._representationInvalid = False if firstHighlightBlock is not None: @@ -633,10 +673,28 @@ class LogPage(Mixins.QTranslatable, QMainWindow): if not isinstance(selectedCategoryId, str): selectedCategoryId = ALL_LOGS_FILTER - sequence, entries = self.manager.snapshot(selectedCategoryId) + cursor = ( + None + if self._representationInvalid + or self._renderedCategoryId != selectedCategoryId + else self._entryCursor + ) + batch = self.manager.entriesSince(cursor, selectedCategoryId) - self._synchronizeDocument(entries, selectedCategoryId) - self._renderedSequence = sequence + try: + self._synchronizeDocument(batch, selectedCategoryId) + except _DocumentBlockAccountingError: + logger.exception( + 'rendered log block accounting diverged; rebuilding document' + ) + + self._entryCursor = None + self._requestRefresh(invalidate=True, immediate=True) + + return + + self._entryCursor = batch.cursor + self._renderedSequence = batch.cursor.sequence self._entriesDirty = False def _scheduleHighlight(self, firstBlock: int): @@ -799,7 +857,7 @@ class LogPage(Mixins.QTranslatable, QMainWindow): @QtCore.Slot(object) def _entriesCleared(self, _categoryIds): """Refresh the presentation after the underlying collection changes.""" - self._requestRefresh(invalidate=True, immediate=True) + self._requestRefresh(immediate=True) def showEvent(self, event): """Render entries accumulated while the page was hidden.""" diff --git a/tests/test_architecture_refactors.py b/tests/test_architecture_refactors.py index 4355937..d59545b 100644 --- a/tests/test_architecture_refactors.py +++ b/tests/test_architecture_refactors.py @@ -1039,6 +1039,40 @@ class ApplicationLifecycleTransactionTest(TestCase): finally: messageQueue.dispose() + def testCoreLogQueueUsesOneNaturalBatchWhenCallbackSupportsIt(self): + """Drain one queue turn through one batch-capable callback invocation.""" + + class BatchCallback: + """Record whether the queue chooses the batch contract.""" + + def __init__(self): + self.single = [] + self.batches = [] + + def __call__(self, message): + self.single.append(message) + + def appendMany(self, messages): + self.batches.append(tuple(messages)) + + callback = BatchCallback() + messageQueue = CoreProcessWorkerModule.MsgQueue(msgCallback=callback) + + try: + messageQueue.getNoWait = mock.Mock( + side_effect=('first', ' ', 'second', 'third', '') + ) + messageQueue.processMsg() + + self.assertEqual(callback.single, []) + self.assertEqual(callback.batches, [('first', 'second', 'third')]) + self.assertEqual( + messageQueue.getTimeout(), + CoreProcessWorkerModule.MsgQueue.ACTIVE_DRAIN_INTERVAL, + ) + finally: + messageQueue.dispose() + def testFailuresAtMeaningfulStagesRollBackOnlyEarlierStages(self): expected = { 'storage': ['cleanup storage', 'cleanup plugins'], diff --git a/tests/test_external_core.py b/tests/test_external_core.py index f2863db..af41e45 100644 --- a/tests/test_external_core.py +++ b/tests/test_external_core.py @@ -140,6 +140,40 @@ class ExternalCoreProcessTest(unittest.TestCase): self.assertFalse(runtime.isAlive()) self.assertFalse(runtime._readerThreads) + def testReaderBatchesCompleteLinesAndPreservesTrailingPartialLine(self): + """Use one callback per read chunk without buffering complete lines.""" + + class BatchCallback: + """Record batch and compatibility callback traffic separately.""" + + def __init__(self): + self.single = [] + self.batches = [] + + def __call__(self, message): + self.single.append(message) + + def appendMany(self, messages): + self.batches.append(tuple(messages)) + + callback = BatchCallback() + runtime = ExternalCoreProcess(msgCallback=callback) + stream = mock.Mock() + stream.fileno.return_value = 123 + + with mock.patch( + 'Furious.Backends.ExternalCore.Process.os.read', + side_effect=(b'first\nsecond\r\npar', b'tial', b''), + ): + runtime._readStream(stream, 'stderr') + + self.assertEqual( + callback.batches, + [('[stderr] first', '[stderr] second')], + ) + self.assertEqual(callback.single, ['[stderr] partial']) + stream.close.assert_called_once_with() + def testMissingExecutableAndImmediateExitFailStartup(self): """Report authoritative path and early non-zero-exit failures.""" with tempfile.TemporaryDirectory(dir=Path.cwd()) as directory: diff --git a/tests/test_log_manager_generation.py b/tests/test_log_manager_generation.py index 6149aca..8b0ddb8 100644 --- a/tests/test_log_manager_generation.py +++ b/tests/test_log_manager_generation.py @@ -30,8 +30,11 @@ from Furious.Service.LogManager import ( CORE_LOG_CATEGORY, TUN2SOCKS_LOG_CATEGORY, LogManager, + formatLogEntry, ) -from Furious.Window.LogPage import LogPage +from Furious.Window.LogPage import LogPage, _documentBlockCount + +from PySide6 import QtGui from collections import OrderedDict from datetime import datetime @@ -320,6 +323,35 @@ class GenerationLogManagerContractTest(unittest.TestCase): _assertManagerInvariants(self, manager, model) return entry + def assertPageMatchesManager(self, page, manager, categoryId=None): + """Prove text, paragraph count, and per-entry ownership agree.""" + entries = manager.entries(categoryId) + reference = QtGui.QTextDocument() + reference.setPlainText('\n'.join(formatLogEntry(entry) for entry in entries)) + + self.assertEqual(page.plainText(), reference.toPlainText()) + self.assertEqual( + page.textBrowser.document().blockCount(), + reference.blockCount(), + ) + self.assertEqual( + tuple( + (metadata.sequence, metadata.blockCount) + for metadata in page._renderedEntries + ), + tuple( + (entry.sequence, _documentBlockCount(formatLogEntry(entry))) + for entry in entries + ), + ) + self.assertEqual( + reference.blockCount(), + max( + 1, + sum(metadata.blockCount for metadata in page._renderedEntries), + ), + ) + def testDeterministicStateTransitionMatrix(self): """Mix every clear kind, rollover, retention mode, and stream.""" manager = self.makeManager( @@ -516,6 +548,394 @@ class GenerationLogManagerContractTest(unittest.TestCase): self.assertEqual(sum(index.oldestRemovals for index in globalIndexes), 1) _assertManagerInvariants(self, manager) + def testIncrementalBatchesReturnOnlyOrderedMissingEntries(self): + """Advance independent global and category cursors over sparse appends.""" + manager = self.makeManager(autoClearEnabled=False) + manager.appendMany( + ('application 1', 'application 2'), + APPLICATION_LOG_CATEGORY, + ) + initial = manager.entriesSince() + coreInitial = manager.entriesSince(None, CORE_LOG_CATEGORY) + + self.assertTrue(initial.resetRequired) + self.assertEqual( + tuple(entry.message for entry in initial.entries), + ('application 1', 'application 2'), + ) + + manager.append('core 1', CORE_LOG_CATEGORY) + manager.append('application 3', APPLICATION_LOG_CATEGORY) + manager.append('core 2', CORE_LOG_CATEGORY) + + globalBatch = manager.entriesSince(initial.cursor) + coreBatch = manager.entriesSince(coreInitial.cursor, CORE_LOG_CATEGORY) + + self.assertFalse(globalBatch.resetRequired) + self.assertEqual( + tuple(entry.message for entry in globalBatch.entries), + ('core 1', 'application 3', 'core 2'), + ) + self.assertEqual( + tuple(entry.message for entry in coreBatch.entries), + ('core 1', 'core 2'), + ) + self.assertEqual( + tuple(entry.sequence for entry in globalBatch.entries), + tuple(sorted(entry.sequence for entry in globalBatch.entries)), + ) + + def testIncrementalBatchReportsRetentionPrefixWithoutFullReset(self): + """Keep a valid cursor while exposing the exact retained prefix boundary.""" + manager = self.makeManager( + maximumEntries=3, + maximumCharacters=1_000, + autoClearEnabled=False, + ) + manager.appendMany(('one', 'two', 'three')) + initial = manager.entriesSince() + + manager.appendMany(('four', 'five')) + batch = manager.entriesSince(initial.cursor) + + self.assertFalse(batch.resetRequired) + self.assertEqual( + tuple(entry.message for entry in batch.entries), + ('four', 'five'), + ) + self.assertEqual(batch.firstRetainedSequence, 3) + self.assertEqual( + tuple(entry.message for entry in manager.entries()), + ('three', 'four', 'five'), + ) + + def testIncrementalCursorsDetectScopedAndSelectiveClears(self): + """Invalidate only views whose structural generation state changed.""" + manager = self.makeManager(autoClearEnabled=False) + manager.registerComponent('other.second', 'Other second', runtime=False) + manager.append('application', APPLICATION_LOG_CATEGORY) + manager.append('core old', CORE_LOG_CATEGORY) + manager.append('other old', 'other.extra') + manager.append('second retained', 'other.second') + + globalCursor = manager.entriesSince().cursor + applicationCursor = manager.entriesSince(None, APPLICATION_LOG_CATEGORY).cursor + coreCursor = manager.entriesSince(None, CORE_LOG_CATEGORY).cursor + + manager.clear(runtimeOnly=True) + manager.append('core new', CORE_LOG_CATEGORY) + + globalBatch = manager.entriesSince(globalCursor) + applicationBatch = manager.entriesSince( + applicationCursor, + APPLICATION_LOG_CATEGORY, + ) + coreBatch = manager.entriesSince(coreCursor, CORE_LOG_CATEGORY) + + self.assertTrue(globalBatch.resetRequired) + self.assertTrue(coreBatch.resetRequired) + self.assertFalse(applicationBatch.resetRequired) + self.assertEqual(applicationBatch.entries, tuple()) + self.assertNotIn( + 'core old', tuple(entry.message for entry in globalBatch.entries) + ) + + globalCursor = globalBatch.cursor + secondCursor = manager.entriesSince(None, 'other.second').cursor + manager.clear('other.extra') + manager.append('second new', 'other.second') + + self.assertTrue(manager.entriesSince(globalCursor).resetRequired) + secondBatch = manager.entriesSince(secondCursor, 'other.second') + self.assertTrue(secondBatch.resetRequired) + self.assertEqual( + tuple(entry.message for entry in secondBatch.entries), + ('second retained', 'second new'), + ) + + def testAppendManyMatchesRepeatedAppendAndCompatibilitySignals(self): + """Preserve normalization, rollover, retention, and signal ordering.""" + timestamp = datetime(2026, 8, 27, 12, 0, 0) + options = { + 'maximumEntries': 7, + 'maximumCharacters': 38, + 'maximumEntryCharacters': 12, + 'autoClearMaximumEntries': 3, + } + batched = self.makeManager(**options) + repeated = self.makeManager(**options) + messages = tuple(f'core-{index}-payload' for index in range(11)) + added = [] + cleared = [] + changed = [] + observedSnapshots = [] + batched.entryAdded.connect(added.append) + batched.entryAdded.connect( + lambda _entry: observedSnapshots.append(batched.entries()) + ) + batched.entriesCleared.connect(cleared.append) + batched.entriesChanged.connect(changed.append) + + batchEntries = batched.appendMany( + messages, + CORE_LOG_CATEGORY, + timestamp=timestamp, + source='batch', + severity='info', + ) + repeatedEntries = tuple( + repeated.append( + message, + CORE_LOG_CATEGORY, + timestamp=timestamp, + source='batch', + severity='info', + ) + for message in messages + ) + + self.assertEqual(batchEntries, repeatedEntries) + self.assertEqual(batched.entries(), repeated.entries()) + self.assertEqual(batched.retainedCharacters, repeated.retainedCharacters) + self.assertEqual(added, list(batchEntries)) + self.assertTrue(observedSnapshots) + self.assertTrue( + all(snapshot == batched.entries() for snapshot in observedSnapshots) + ) + self.assertEqual(len(cleared), 3) + processQtEvents() + self.assertEqual(changed, [batchEntries[-1].sequence]) + self.assertLessEqual(batched.entryCount(), batched.maximumEntries) + self.assertLessEqual(batched.retainedCharacters, batched.maximumCharacters) + _assertManagerInvariants(self, batched) + _assertManagerInvariants(self, repeated) + + def testProducerCallbackSupportsSingleAndAtomicBatchDelivery(self): + """Return stored entries through both safe producer callback paths.""" + manager = self.makeManager(autoClearEnabled=False) + callback = manager.callback( + CORE_LOG_CATEGORY, + source='producer', + severity='info', + ) + + batched = callback.appendMany(('one', 'two', 'three')) + single = callback('four') + + self.assertEqual( + tuple(entry.message for entry in manager.entries(CORE_LOG_CATEGORY)), + ('one', 'two', 'three', 'four'), + ) + self.assertEqual( + tuple(entry.message for entry in batched), ('one', 'two', 'three') + ) + self.assertIs(single, manager.entries(CORE_LOG_CATEGORY)[-1]) + self.assertTrue( + all( + entry.source == 'producer' and entry.severity == 'info' + for entry in manager.entries(CORE_LOG_CATEGORY) + ) + ) + + def testMultilineEntryEvictionRemovesItsExactDocumentBlocks(self): + """Never let paragraph retention silently retain a fragment of an entry.""" + with isolatedSettings(): + manager = self.makeManager( + maximumEntries=2, + maximumCharacters=10_000, + maximumEntryCharacters=1_000, + autoClearEnabled=False, + ) + page = LogPage(manager=manager) + manager.append('a\nb\nc\nd\ne') + page.show() + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager) + + renderedText = page.plainText() + renderedMetadata = tuple(page._renderedEntries) + + with self.assertRaises(RuntimeError): + page._removeLeadingDocumentBlocks( + page.textBrowser.document().blockCount() + 1 + ) + + self.assertEqual(page.plainText(), renderedText) + self.assertEqual(tuple(page._renderedEntries), renderedMetadata) + + manager.append('x') + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager) + + firstMetadata = page._renderedEntries[0] + page._renderedEntries[0] = type(firstMetadata)( + sequence=firstMetadata.sequence, + blockCount=page.textBrowser.document().blockCount() + 1, + ) + + with self.assertLogs('Furious.Window.LogPage', level='ERROR') as logs: + manager.append('y') + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + + self.assertTrue( + any('rebuilding document' in message for message in logs.output) + ) + self.assertPageMatchesManager(page, manager) + + self.assertEqual(page.plainText(), 'x\ny') + self.assertEqual(page.textBrowser.document().blockCount(), 2) + + page.close() + page.deleteLater() + collectAtBoundary() + + def testMixedMultilineRetentionFilterAndClearStaySynchronized(self): + """Keep exact ownership through mixed separators, truncation, and resets.""" + with isolatedSettings(): + manager = self.makeManager( + maximumEntries=5, + maximumCharacters=10_000, + maximumEntryCharacters=48, + autoClearEnabled=False, + ) + page = LogPage(manager=manager) + page.show() + + messages = ( + ('alpha\r\nbeta', APPLICATION_LOG_CATEGORY), + ('gamma\rdelta', CORE_LOG_CATEGORY), + ('epsilon\u2029zeta', APPLICATION_LOG_CATEGORY), + ('eta\u2028theta', CORE_LOG_CATEGORY), + ('iota\n\nkappa', APPLICATION_LOG_CATEGORY), + ('x' * 80, CORE_LOG_CATEGORY), + ('', APPLICATION_LOG_CATEGORY), + ('lambda\r\nmu\rnu\u2029xi', CORE_LOG_CATEGORY), + ) + + for message, categoryId in messages: + manager.append(message, categoryId) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager) + + coreIndex = page.filterComboBox.findData(CORE_LOG_CATEGORY) + page.filterComboBox.setCurrentIndex(coreIndex) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager, CORE_LOG_CATEGORY) + + manager.clear(CORE_LOG_CATEGORY) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager, CORE_LOG_CATEGORY) + + manager.append('runtime\ncore', CORE_LOG_CATEGORY) + manager.append('runtime\rcomponent', TUN2SOCKS_LOG_CATEGORY) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager, CORE_LOG_CATEGORY) + + manager.clear(runtimeOnly=True) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager, CORE_LOG_CATEGORY) + + allIndex = page.filterComboBox.findData('all') + page.filterComboBox.setCurrentIndex(allIndex) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager) + + manager.append('before\nfull clear') + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + manager.clear() + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertPageMatchesManager(page, manager) + + page.close() + page.deleteLater() + collectAtBoundary() + + def testMultilineRetentionKeepsProgressiveHighlightingAligned(self): + """Shift pending highlighting by removed blocks rather than entries.""" + with isolatedSettings(): + manager = self.makeManager( + maximumEntries=3, + maximumCharacters=20_000, + maximumEntryCharacters=4_000, + autoClearEnabled=False, + ) + page = LogPage(manager=manager) + + def message(index): + return '\n'.join( + '2026/08/27 18:45:' + f'{second:02d}.000000 from 127.0.0.1:' + f'{5000 + index * 10 + second} accepted ' + '//example.com:443 [http >> proxy]' + for second in range(3) + ) + + manager.appendMany(tuple(message(index) for index in range(3))) + page.show() + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertTrue(waitFor(lambda: page._highlightNextBlock is None)) + + manager.appendMany((message(3), message(4))) + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertTrue(waitFor(lambda: page._highlightNextBlock is None)) + self.assertPageMatchesManager(page, manager) + + document = page.textBrowser.document() + missingFormats = tuple( + blockNumber + for blockNumber in range(document.blockCount()) + if not document.findBlockByNumber(blockNumber).layout().formats() + ) + self.assertEqual(missingFormats, tuple()) + + page.close() + page.deleteLater() + collectAtBoundary() + + def testLogPagePullsOneSuffixInsteadOfRepeatedFullSnapshots(self): + """Keep a visible page on the cursor path during a high-frequency burst.""" + with isolatedSettings(): + manager = self.makeManager( + maximumEntries=500, + maximumCharacters=20_000, + autoClearEnabled=False, + ) + manager.appendMany(tuple(f'initial {index}' for index in range(200))) + fullReads = [] + suffixReads = [] + originalFullRead = manager._entriesLocked + originalSuffixRead = manager._entriesAfterLocked + + def fullRead(categoryId): + fullReads.append(categoryId) + return originalFullRead(categoryId) + + def suffixRead(sequence, categoryId): + suffixReads.append((sequence, categoryId)) + return originalSuffixRead(sequence, categoryId) + + manager._entriesLocked = fullRead + manager._entriesAfterLocked = suffixRead + page = LogPage(manager=manager) + page.resize(900, 420) + page.show() + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertEqual(len(fullReads), 1) + + for index in range(250): + manager.append(f'burst {index}', APPLICATION_LOG_CATEGORY) + + self.assertTrue(waitFor(lambda: not page._entriesDirty)) + self.assertEqual(len(fullReads), 1) + self.assertEqual(len(suffixReads), 1) + self.assertEqual( + page.plainText().splitlines(), + [entry.message for entry in manager.entries()], + ) + page.close() + page.deleteLater() + collectAtBoundary() + 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)): @@ -704,6 +1124,8 @@ class GenerationLogManagerContractTest(unittest.TestCase): before = manager.entries() with self.assertRaises(RuntimeError): manager.append(RaisingString()) + with self.assertRaises(RuntimeError): + manager.appendMany(('must not commit', RaisingString(), 'unreached')) with self.assertRaises(TypeError): manager.append('bad timestamp', timestamp='not a datetime') with self.assertRaises(KeyError): @@ -835,12 +1257,28 @@ class GenerationLogManagerContractTest(unittest.TestCase): def reader(): try: start.wait(5) + cursor = None 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)), ) + batch = manager.entriesSince(cursor, CORE_LOG_CATEGORY) + self.assertEqual( + tuple(item.sequence for item in batch.entries), + tuple(sorted(item.sequence for item in batch.entries)), + ) + + if cursor is not None and not batch.resetRequired: + self.assertTrue( + all( + item.sequence > cursor.sequence + for item in batch.entries + ) + ) + + cursor = batch.cursor manager.entryCount(CORE_LOG_CATEGORY) except Exception as error: errors.append(error) @@ -1169,7 +1607,7 @@ class VeryHeavyGenerationLogManagerTest(unittest.TestCase): ) def testCleanupLatencyDoesNotScaleWithBacklog(self): - """Keep append cleanup capped at 64 for geometrically larger queues.""" + """Keep append cleanup capped at one for geometrically larger queues.""" results = [] for backlog in (64, 1_000, 10_000): samples = [] @@ -1181,6 +1619,7 @@ class VeryHeavyGenerationLogManagerTest(unittest.TestCase): autoClearEnabled=False, ) manager.RetiredCleanupBudget = 64 + manager.AppendRetiredCleanupBudget = 1 for index in range(backlog): manager.append(str(index), CORE_LOG_CATEGORY) manager.clear(runtimeOnly=True) @@ -1188,7 +1627,7 @@ class VeryHeavyGenerationLogManagerTest(unittest.TestCase): started = time.perf_counter_ns() manager.append('new') samples.append(time.perf_counter_ns() - started) - self.assertEqual(manager.retiredEntryCount, max(0, before - 64)) + self.assertEqual(manager.retiredEntryCount, max(0, before - 1)) 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) diff --git a/tests/test_models_and_services.py b/tests/test_models_and_services.py index 80a4acb..159ad9a 100644 --- a/tests/test_models_and_services.py +++ b/tests/test_models_and_services.py @@ -815,6 +815,7 @@ class LogManagerTest(unittest.TestCase): """Release at most the fixed budget from a retired generation per turn.""" manager = LogManager(maximumEntries=100, autoClearEnabled=False) manager.RetiredCleanupBudget = 3 + manager.AppendRetiredCleanupBudget = 1 entryReferences = [] for index in range(10): @@ -835,8 +836,8 @@ class LogManagerTest(unittest.TestCase): manager.append('application', APPLICATION_LOG_CATEGORY) - self.assertEqual(manager.retiredEntryCount, 7) - self.assertEqual(observedEntries.oldestRemovals, 3) + self.assertEqual(manager.retiredEntryCount, 9) + self.assertEqual(observedEntries.oldestRemovals, 1) cleanupTurns = 0