mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 22:35:08 +00:00
* fix(errors): record the exception's own stack trace record_error() called traceback.format_exc(), which only sees an exception while its except block is running. plugin_executor records exceptions caught on a worker thread after that block has ended, so every trace on /errors read "NoneType: None". The trace is now built from the exception's __traceback__. The executor's log call had the same problem with exc_info=True and now passes the exception. record_error() also merged LEDMatrixError context into the caller's dict in place; it now works on a copy. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(wifi): point at configure_wifi_permissions.sh instead of a sudoers list The module docstring told users to grant NOPASSWD sudo on iptables and ip. configure_wifi_permissions.sh refuses those grants on purpose: a wildcard rule for either runs an arbitrary program as root. Point at the script and say why it leaves them out. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(wifi): disconnect finds the saved profile by SSID disconnect_from_network() asked `nmcli -f NAME,802-11-wireless.ssid connection show` for the profile to take down, but nmcli rejects that column for `connection show`, so the lookup always failed and only the device was disconnected. The per-profile lookup _connect_nmcli() already used is now _find_profile_for_ssid(), and both callers share it. It also splits terse output on the last colon and unescapes "\:", so a profile name containing a colon is found. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(wifi): write wifi_config.json atomically and report a failed save _save_config() opened the file for writing in place and swallowed any error, so a wifi_config.json left owned by root made the web toggle for auto-enabling AP mode report success while nothing was saved, and a crash mid-write could truncate the file. It now uses atomic_write_json, which also keeps the file's owner and shared group when root saves it, and returns False on failure. POST /wifi/ap/auto-enable answers 500 in that case. The file is now written with indent=4, like the other config files. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(fonts): resolve plugin:// fonts in the plugin's own directory FontManager looked for a plugin's bundled fonts under Path("plugins") / plugin_id: relative to the process cwd, and not the default install directory (plugin-repos/), so a manifest's plugin:// fonts never loaded. register_plugin_fonts() takes an optional plugin_dir, and PluginManager passes the directory it loaded the plugin from. Callers that omit it get a lookup in the configured plugin_system.plugins_directory, then plugins/, resolved against the install root. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(api-helper): cache responses for the requested cache_ttl APIHelper.get(cache_ttl=...) and set_cache(ttl=...) dropped the ttl on the claim that CacheManager does not support one, but CacheManager.set() takes a ttl, stores it with the entry, and both cache tiers honour it over a reader's max_age. Without it every response expired after the 300-second default read age, whatever the plugin asked for. The ttl is now passed through, and the cache read passes cache_ttl as max_age for entries written without one. The class docstring describes what the helper actually does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(style): one scale range for the schema, element_scale and LogoHelper The generated Scale field allowed 0.1 to 10, element_style's reader capped at 10 with no floor, and LogoHelper accepted 0.05 to 8 and reset anything else to 1.0. A logo scale of 9, which the form accepts, drew at the shipped size. MIN_ELEMENT_SCALE / MAX_ELEMENT_SCALE (0.1, 10.0) in src.element_style are now the schema bounds and the clamp every reader applies through coerce_scale(): a positive number outside the range is clamped, and anything that is not a finite positive number means the default. That also stops element_scale() passing NaN through, since min(nan, 10.0) is nan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(logos): placeholder lands at the requested path; empty logos list download_missing_logo() wrote its fallback placeholder to <normalize_abbreviation(abbr)>.png in the logo directory rather than to the logo_path the caller passed, so it could return True while nothing existed where the plugin looks (e.g. "TA&M.png" vs "TAANDM.png"). create_placeholder_logo() takes an optional filepath, and download_missing_logo passes the requested one. download_missing_logo_for_team() only caught KeyError, so a team whose "logos" list is empty raised IndexError; it now treats KeyError, IndexError and TypeError as "no logo URL". The placeholder is drawn with PLACEHOLDER_SIZE / PLACEHOLDER_BG, the constants is_placeholder_logo() recognises it by, instead of repeated literals. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(fonts): resolve bundled font paths against the install root TextHelper's default font_dir, the logo placeholder's font and FontManager's font_overrides.json were all relative to the process cwd, so a process started anywhere but the install root (the plugin safety harness, a manual run, a unit without WorkingDirectory) drew with PIL's default face and read no overrides. They now go through font_layout.resolve_asset_path; the overrides file sits in the install root's config/. The resolver docstrings described an order the code does not follow: resolve_asset_path never consults the cwd, and sports_shared's _resolve_font_path tries the cwd first. Both docstrings now say what the code does, and _resolve_font_path calls resolve_asset_path instead of probing FontManager for it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(sync): the web UI reads the sync status file the display writes sync_manager writes its status to tempfile.gettempdir(), but GET /api/v3/sync/status read a hardcoded /tmp/led_matrix_sync_status.json and defaulted the port to a literal 5765. Wherever TMPDIR is set (or on any non-/tmp host) the page only ever showed "starting". The endpoint now uses sync_manager.STATUS_FILE and SYNC_PORT. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(http): the rankings resolver sends the project's User-Agent DynamicTeamResolver fetched ESPN rankings with a bare requests.get, so it sent python-requests' default User-Agent, which ESPN rejects; the AP_TOP_N favourites then resolved to nothing. It now sends DEFAULT_HTTP_HEADERS. BaseOddsManager carried its own copy of the User-Agent string and now uses the same shared headers (which also adds Accept-Language). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(backup): record the core release and read the configured plugin dir The manifest's ledmatrix_version came from a VERSION file that does not exist, then from .git/HEAD: a 12-character sha, or "ref: refs/he" when the branch's ref was packed. It is now src.__version__. list_installed_plugins() scanned a hardcoded plugin-repos/, so on an install whose plugin_system.plugins_directory points elsewhere, plugins missing from plugin_state.json were left out of the backup. It now reads the configured directory from config/config.json, defaulting to plugin-repos. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(startup): report a missing display section once A config without a display section produced three errors for the one problem ("Missing required configuration key: display", "Display configuration is missing or empty" and "Display configuration is missing"), and an empty one produced two. _validate_config now reports it once, as a missing key or an empty section, and _validate_display_config leaves it to that. The module docstring said the validator fails fast; nothing in the display service calls raise_on_errors(), so it now says the errors are reported and startup continues. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(wifi): share the copied blocks and name the AP constants - _parse_nmcli_wifi_list() is the one parser behind _scan_nmcli and _scan_nmcli_cached. - _verify_connected(), _wait_for_device_idle(), _failsafe_ap() and _mark_forced() replace blocks that were pasted two or three times in the connect and enable-AP paths. The device-idle wait now checks before its first one-second sleep instead of after it. - _check_command() calls _find_command_path() instead of repeating it. - AP_IP, PORTAL_PORT, AP_PROFILE_NAME and AP_PROFILE_NAMES name values that were spelled out 14, 12, 8 and 2 times; the two deletion loops now walk the same tuple. The iwconfig status path compares the AP address exactly: startswith() also skipped 192.168.4.10-19. - Dropped a second WIFI.SIGNAL query that repeated the first, a no-op "if ssid: continue", the try/except around _connect_wpa_supplicant's constant return, and a second save of a scan scan_networks already saves. - _ensure_wifi_radio_enabled's docstring says it returns True when the radio state cannot be read at all. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(config): drop dead branches and history comments in ConfigManager - The module docstring pointed plugin authors at update_plugin_config(), which does not exist; it now names save_config_atomic() and save_raw_file_content(). - load_config's FileNotFoundError handler tested the message for "config_secrets.json", but a missing secrets file is handled where it is read, so only config.json reaches it; the check is gone. - save_raw_file_content's `file_type == "main" or "secrets"` guard was always true (anything else raised earlier). - get_raw_file_content('secrets') already returns {} for a missing file, so the os.path.exists() in front of two calls to it is gone. - Comments that narrated earlier behaviour are rewritten as what the code does now. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(background-data): present-tense comments, drop unused API - Comments that told the history of each fix (what "used to" happen, "the old per-delivery release") now state the invariant the code keeps. - get_statistics() no longer reports a constant 'queue_size': 0, and the uncalled clear_completed_requests() is gone (_cleanup_completed_requests does that job on every completion). Neither is referenced in core, the web UI or the plugin monorepo. shutdown_background_service() has no production caller either, but it is the only way to tear down the get_background_service() singleton, which the tests rely on, so it stays. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(odds): drop the unread cache_ttl and merge the odds_data branches BaseOddsManager loaded base_odds_manager.cache_ttl from config and never used it: cached odds live for the update interval (get_odds' ttl=interval). No core or monorepo code reads the attribute, so it is gone along with its log line. The two consecutive `if odds_data:` blocks are one. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(backup): one table for the single-file sections config, secrets, wifi and ytm_auth were each spelled out in create, preview, validate and restore. _SINGLE_FILE_SECTIONS lists them once, with the RestoreOptions flag that restores each, and all four walk it. Restore error messages keep their wording ("Failed to restore <file name>"). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(fonts): drop FontManager's write-only state and duplicate logs - fonts_config, font_metadata and font_dependencies were written and never read; the performance_stats keys font_load_times, render_times, total_renders and the per-call "resolve" timings (_record_performance_metric) likewise. get_performance_stats() reads only the counters that remain. Nothing in core or the plugin monorepo references any of them. - A failed BDF load was logged twice, by _load_bdf_font and again by get_font; get_font's line is the one kept. - Removed "NEW:" and commented-out cozette entries, the "Copy font to assets/fonts" comment on code that copies nothing, and local imports of names the module already imports. The deprecated add_font() now resolves assets/fonts against the install root. The @deprecated methods stay. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(text-helper): cache loaded fonts; drop the pre-textlength fallback TextHelper declared _font_cache, cleared it and reported its size, but never stored anything in it. load_fonts() now keeps each (file, size) it loads there, so clear_font_cache() and get_font_cache_stats() mean what they say and repeated load_fonts() calls reuse the fonts. get_text_width() no longer catches AttributeError for Pillow releases without ImageDraw.textlength; requirements.txt pins Pillow>=12.2. The class docstring describes what the helper does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(common): fix wrong docstrings in api_helper, permission_utils, snapshot_policy - permission_utils called 0o2775 "sticky bit"; the 2 is setgid, which is what makes new files take the directory's group. - snapshot_policy pointed at web_interface/blueprints/api_v3.py, which is a package now; the health check is in api_v3/misc.py. - APIHelper.clear_cache() lost a history note and a fallback to a clear() method that neither CacheManager nor the testing MockCacheManager has. The session headers are built from DEFAULT_HTTP_HEADERS instead of a copy of them, and the module docstring says what the module offers. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(sports): present-tense comments in the shared scoreboard renderers - sports_scroll and sports_game_renderer comments that referred to "this PR", "the old flat 128px card" or what the renderer "previously" did now describe the current behaviour and its reason. - The block explaining why non-finite settings are rejected sat above _score_reserve_width; it describes _center_gap_width and now lives in it. - unshare_element_fonts wrapped its import of font_layout.load_truetype in an `except ImportError` that cannot fire inside core; the import stays at call time so tests can spy on the pinned loader. - sports_card docstrings that told the history of a fix say what the code does. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(sports-shared): drop dead code, name the ESPN limit - _get_weeks_data asked for limit=1000, which fetch_espn_scoreboard clamps to ESPN_MAX_LIMIT anyway; it now names that constant. Its unused `immediate_events = []` is gone. - _get_season_schedule_dates() returned ("", "") and has no caller in core or the plugin monorepo. - _should_log keeps its warning_type parameter (part of the inherited signature, though nothing in core or the monorepo calls it) and its docstring says the cooldown is shared across types. - An unused ImageFont import is gone. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(sync): one follower-mode switch, shared panel defaults - The class docstring said the leader sends PNG frames. Frames go over UDP as raw RGB; PNG is only the Vegas scroll image sent over TCP. It now describes both paths. - _enter_follower_mode() replaces the two copies of "note the leader, switch from standalone to follower, log, write status" in the frame and scroll-position handlers. - The rows/cols fallbacks use DEFAULT_ROWS / DEFAULT_COLS from src.display_geometry, as chain_length already did. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(style): drop _layout_axis, name the layout group title - ElementStyleResolver._layout_axis() had no caller in core or the plugin monorepo. - _element_block_from_spec checked spec['size'] was a dict again after size_spec already had; it reads size_spec. - The "Layout Offsets" title written into three generated schema blocks is _LAYOUT_TITLE. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(logo-helper): say what the placeholder draws; name the 1.5 box factor - _create_placeholder_logo's docstring said it draws the team abbreviation; it draws an outlined grey box and nothing else. The docstring says so, and the "in a real implementation you'd want text" comments are gone. - The 1.5 x panel default logo box, written out six times, is DEFAULT_LOGO_BOX_FACTOR. - ImageDraw is imported with Image at the top of the module. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(logos): drop dead code and a duplicate regex in logo_downloader - _SAFE_LEAGUE_CODE_RE was the same pattern as _SAFE_LEAGUE_RE; both checks use the one. - get_logo_filename_variations reassigned the TA&M case to the list it already had; the function returns the two names directly. - _get_team_name_variations() had no caller in core or the plugin monorepo. - fetch_single_team's docstring was copied from fetch_teams_data; a log message read "for{team_id}". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor: drop the Pillow<9.1 resample shim and a catch-and-reraise - adaptive_images fell back to Image.LANCZOS/NEAREST for Pillow < 9.1; requirements.txt pins Pillow>=12.2. RESAMPLE_LANCZOS and RESAMPLE_NEAREST keep their names (src.common re-exports them). - CacheManager.save_cache caught CacheError only to re-raise it; the disk write is now called directly, with the same result. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * test(api-helper): stop the real CacheManager's cleanup thread The cache-lifetime tests built a CacheManager and left its cleanup thread's class-wide claim on the directory in place, which broke test_cache_cleanup_thread_ownership when it ran later in the session. The fixture now stops the thread on teardown. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(changelog): core-common Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
903 lines
36 KiB
Python
903 lines
36 KiB
Python
"""
|
|
Error Aggregation Service
|
|
|
|
Provides centralized error tracking, pattern detection, and reporting
|
|
for the LEDMatrix system. Enables automatic bug detection by tracking
|
|
error frequency, patterns, and context.
|
|
|
|
This is a local-only implementation with no external dependencies.
|
|
Errors are stored in memory with optional JSON export.
|
|
"""
|
|
|
|
import math
|
|
import threading
|
|
import time
|
|
import traceback
|
|
import json
|
|
import uuid
|
|
from collections import defaultdict
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timedelta
|
|
from pathlib import Path
|
|
from typing import Dict, List, Optional, Any, Callable, Tuple
|
|
import logging
|
|
|
|
from src.exceptions import LEDMatrixError
|
|
from src.redaction import redact_credentials
|
|
|
|
|
|
def _format_trace(error: BaseException) -> str:
|
|
"""The traceback carried by ``error`` itself.
|
|
|
|
Callers often record an exception after its ``except`` block has ended,
|
|
or from another thread than the one that raised it (plugin_executor runs
|
|
plugins on worker threads), where ``traceback.format_exc()`` has nothing
|
|
to report. The exception object keeps its own ``__traceback__``, so the
|
|
trace is built from that. An exception that was created but never raised
|
|
has no traceback, and the result is just its type and message.
|
|
"""
|
|
return "".join(traceback.format_exception(type(error), error, error.__traceback__))
|
|
|
|
|
|
@dataclass
|
|
class ErrorRecord:
|
|
"""Record of a single error occurrence."""
|
|
error_type: str
|
|
message: str
|
|
timestamp: datetime
|
|
context: Dict[str, Any] = field(default_factory=dict)
|
|
plugin_id: Optional[str] = None
|
|
operation: Optional[str] = None
|
|
stack_trace: Optional[str] = None
|
|
|
|
def to_dict(self) -> Dict[str, Any]:
|
|
"""Convert to dictionary for JSON serialization."""
|
|
return {
|
|
"error_type": self.error_type,
|
|
"message": self.message,
|
|
"timestamp": self.timestamp.isoformat(),
|
|
"context": self.context,
|
|
"plugin_id": self.plugin_id,
|
|
"operation": self.operation,
|
|
"stack_trace": self.stack_trace
|
|
}
|
|
|
|
|
|
@dataclass
|
|
class ErrorPattern:
|
|
"""Detected error pattern for automatic detection."""
|
|
error_type: str
|
|
count: int
|
|
first_seen: datetime
|
|
last_seen: datetime
|
|
affected_plugins: List[str] = field(default_factory=list)
|
|
sample_messages: List[str] = field(default_factory=list)
|
|
severity: str = "warning" # warning, error, critical
|
|
|
|
def to_dict(self) -> Dict[str, Any]:
|
|
"""Convert to dictionary for JSON serialization."""
|
|
return {
|
|
"error_type": self.error_type,
|
|
"count": self.count,
|
|
"first_seen": self.first_seen.isoformat(),
|
|
"last_seen": self.last_seen.isoformat(),
|
|
"affected_plugins": list(dict.fromkeys(self.affected_plugins)),
|
|
"sample_messages": self.sample_messages[:3], # Keep only 3 samples
|
|
"severity": self.severity
|
|
}
|
|
|
|
|
|
class ErrorAggregator:
|
|
"""
|
|
Aggregates and analyzes errors across the system.
|
|
|
|
Features:
|
|
- Error counting by type, plugin, and time window
|
|
- Pattern detection (recurring errors)
|
|
- Error rate alerting via callbacks
|
|
- Export for analytics/reporting
|
|
|
|
Thread-safe for concurrent access.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
max_records: int = 1000,
|
|
pattern_threshold: int = 5,
|
|
pattern_window_minutes: int = 60,
|
|
export_path: Optional[Path] = None
|
|
):
|
|
"""
|
|
Initialize the error aggregator.
|
|
|
|
Args:
|
|
max_records: Maximum number of error records to keep in memory
|
|
pattern_threshold: Number of occurrences to detect a pattern
|
|
pattern_window_minutes: Time window for pattern detection
|
|
export_path: Optional path for JSON export (auto-export on pattern detection)
|
|
"""
|
|
self.logger = logging.getLogger(__name__)
|
|
self.max_records = max_records
|
|
self.pattern_threshold = pattern_threshold
|
|
self.pattern_window = timedelta(minutes=pattern_window_minutes)
|
|
self.export_path = export_path
|
|
|
|
self._records: List[ErrorRecord] = []
|
|
self._error_counts: Dict[str, int] = defaultdict(int)
|
|
self._plugin_error_counts: Dict[str, Dict[str, int]] = defaultdict(lambda: defaultdict(int))
|
|
self._patterns: Dict[str, ErrorPattern] = {}
|
|
self._pattern_callbacks: List[Callable[[ErrorPattern], None]] = []
|
|
self._lock = threading.RLock() # RLock allows nested acquisition for export_to_file
|
|
|
|
# Track session start for relative timing
|
|
self._session_start = datetime.now()
|
|
|
|
# Bumped on every change, so a publisher can tell "nothing new since
|
|
# the last snapshot" without comparing snapshots.
|
|
self._version = 0
|
|
|
|
def record_error(
|
|
self,
|
|
error: Exception,
|
|
context: Optional[Dict[str, Any]] = None,
|
|
plugin_id: Optional[str] = None,
|
|
operation: Optional[str] = None
|
|
) -> ErrorRecord:
|
|
"""
|
|
Record an error occurrence.
|
|
|
|
Args:
|
|
error: The exception that occurred
|
|
context: Optional context dictionary with additional details
|
|
plugin_id: Optional plugin ID that caused the error
|
|
operation: Optional operation name (e.g., "update", "display")
|
|
|
|
Returns:
|
|
The created ErrorRecord
|
|
"""
|
|
with self._lock:
|
|
error_type = type(error).__name__
|
|
|
|
# A copy, so the caller's dict is not changed behind its back.
|
|
error_context = dict(context) if context else {}
|
|
if isinstance(error, LEDMatrixError) and error.context:
|
|
error_context.update(error.context)
|
|
|
|
record = ErrorRecord(
|
|
error_type=error_type,
|
|
message=str(error),
|
|
timestamp=datetime.now(),
|
|
context=error_context,
|
|
plugin_id=plugin_id,
|
|
operation=operation,
|
|
stack_trace=_format_trace(error)
|
|
)
|
|
|
|
# Add record (with size limit)
|
|
self._records.append(record)
|
|
if len(self._records) > self.max_records:
|
|
self._records.pop(0)
|
|
|
|
# Update counts
|
|
self._error_counts[error_type] += 1
|
|
if plugin_id:
|
|
self._plugin_error_counts[plugin_id][error_type] += 1
|
|
self._version += 1
|
|
|
|
# Check for patterns
|
|
self._detect_pattern(record)
|
|
|
|
# Log the error
|
|
self.logger.debug(
|
|
f"Error recorded: {error_type} - {str(error)[:100]}",
|
|
extra={"plugin_id": plugin_id, "operation": operation}
|
|
)
|
|
|
|
return record
|
|
|
|
def _detect_pattern(self, record: ErrorRecord) -> None:
|
|
"""Detect recurring error patterns."""
|
|
cutoff = datetime.now() - self.pattern_window
|
|
recent_same_type = [
|
|
r for r in self._records
|
|
if r.error_type == record.error_type and r.timestamp > cutoff
|
|
]
|
|
|
|
if len(recent_same_type) >= self.pattern_threshold:
|
|
pattern_key = record.error_type
|
|
is_new_pattern = pattern_key not in self._patterns
|
|
|
|
# Determine severity based on count
|
|
count = len(recent_same_type)
|
|
if count > self.pattern_threshold * 3:
|
|
severity = "critical"
|
|
elif count > self.pattern_threshold * 2:
|
|
severity = "error"
|
|
else:
|
|
severity = "warning"
|
|
|
|
# Collect affected plugins
|
|
# Unique, in first-seen order. Each repeat re-scans the whole
|
|
# window, so merging duplicates in below grew without bound.
|
|
affected_plugins = list(dict.fromkeys(r.plugin_id for r in recent_same_type if r.plugin_id))
|
|
|
|
# Collect sample messages
|
|
sample_messages = list(set(r.message for r in recent_same_type[:5]))
|
|
|
|
if is_new_pattern:
|
|
pattern = ErrorPattern(
|
|
error_type=record.error_type,
|
|
count=count,
|
|
first_seen=recent_same_type[0].timestamp,
|
|
last_seen=record.timestamp,
|
|
affected_plugins=affected_plugins,
|
|
sample_messages=sample_messages,
|
|
severity=severity
|
|
)
|
|
self._patterns[pattern_key] = pattern
|
|
|
|
self.logger.warning(
|
|
f"Error pattern detected: {record.error_type} occurred "
|
|
f"{count} times in last {self.pattern_window}. "
|
|
f"Affected plugins: {set(affected_plugins) or 'unknown'}"
|
|
)
|
|
|
|
# Notify callbacks
|
|
for callback in self._pattern_callbacks:
|
|
try:
|
|
callback(pattern)
|
|
except Exception as e:
|
|
self.logger.error(f"Pattern callback failed: {e}")
|
|
|
|
# Auto-export if path configured
|
|
if self.export_path:
|
|
self._auto_export()
|
|
else:
|
|
# Update existing pattern
|
|
self._patterns[pattern_key].count = count
|
|
self._patterns[pattern_key].last_seen = record.timestamp
|
|
self._patterns[pattern_key].severity = severity
|
|
known = self._patterns[pattern_key].affected_plugins
|
|
known.extend(p for p in affected_plugins if p not in known)
|
|
|
|
def on_pattern_detected(self, callback: Callable[[ErrorPattern], None]) -> None:
|
|
"""
|
|
Register a callback to be called when a new error pattern is detected.
|
|
|
|
Args:
|
|
callback: Function that takes an ErrorPattern as argument
|
|
"""
|
|
self._pattern_callbacks.append(callback)
|
|
|
|
def get_error_summary(self) -> Dict[str, Any]:
|
|
"""
|
|
Get summary of all errors for reporting.
|
|
|
|
Returns:
|
|
Dictionary with error statistics and recent errors
|
|
"""
|
|
with self._lock:
|
|
# Calculate error rate (errors per hour)
|
|
session_duration = (datetime.now() - self._session_start).total_seconds() / 3600
|
|
error_rate = len(self._records) / max(session_duration, 0.01)
|
|
|
|
return {
|
|
"session_start": self._session_start.isoformat(),
|
|
"total_errors": len(self._records),
|
|
"error_rate_per_hour": round(error_rate, 2),
|
|
"error_counts_by_type": dict(self._error_counts),
|
|
"plugin_error_counts": {
|
|
k: dict(v) for k, v in self._plugin_error_counts.items()
|
|
},
|
|
"active_patterns": {
|
|
k: v.to_dict() for k, v in self._patterns.items()
|
|
},
|
|
"recent_errors": [
|
|
r.to_dict() for r in self._records[-20:]
|
|
]
|
|
}
|
|
|
|
def get_plugin_health(self, plugin_id: str) -> Dict[str, Any]:
|
|
"""
|
|
Get health status for a specific plugin.
|
|
|
|
Args:
|
|
plugin_id: Plugin ID to check
|
|
|
|
Returns:
|
|
Dictionary with plugin error statistics
|
|
"""
|
|
with self._lock:
|
|
plugin_errors = self._plugin_error_counts.get(plugin_id, {})
|
|
recent_plugin_errors = [
|
|
r for r in self._records[-100:]
|
|
if r.plugin_id == plugin_id
|
|
]
|
|
|
|
# Determine health status
|
|
recent_count = len(recent_plugin_errors)
|
|
if recent_count == 0:
|
|
status = "healthy"
|
|
elif recent_count < 5:
|
|
status = "degraded"
|
|
else:
|
|
status = "unhealthy"
|
|
|
|
return {
|
|
"plugin_id": plugin_id,
|
|
"status": status,
|
|
"total_errors": sum(plugin_errors.values()),
|
|
"error_types": dict(plugin_errors),
|
|
"recent_error_count": recent_count,
|
|
"last_error": recent_plugin_errors[-1].to_dict() if recent_plugin_errors else None
|
|
}
|
|
|
|
def clear_old_records(self, max_age_hours: int = 24) -> int:
|
|
"""
|
|
Clear records older than specified age.
|
|
|
|
Args:
|
|
max_age_hours: Maximum age in hours
|
|
|
|
Returns:
|
|
Number of records cleared
|
|
"""
|
|
with self._lock:
|
|
cutoff = datetime.now() - timedelta(hours=max_age_hours)
|
|
original_count = len(self._records)
|
|
self._records = [r for r in self._records if r.timestamp > cutoff]
|
|
cleared = original_count - len(self._records)
|
|
|
|
if cleared > 0:
|
|
self.logger.info(f"Cleared {cleared} old error records")
|
|
|
|
return cleared
|
|
|
|
@property
|
|
def version(self) -> int:
|
|
"""Changes whenever the recorded errors do (see ErrorSnapshotPublisher)."""
|
|
return self._version
|
|
|
|
def clear_before(self, cutoff: datetime) -> int:
|
|
"""Forget every error recorded at or before ``cutoff``.
|
|
|
|
Unlike clear_old_records, this also resets what the summary reports:
|
|
the per-type and per-plugin counts are rebuilt from the records that
|
|
remain, and detected patterns that began before the cutoff are dropped
|
|
(one that is still happening is detected again on its next
|
|
occurrence). Errors recorded after the
|
|
cutoff are kept, so a clear that is applied a few seconds after it was
|
|
requested does not swallow what happened in between.
|
|
|
|
Returns:
|
|
Number of records removed
|
|
"""
|
|
with self._lock:
|
|
kept = [r for r in self._records if r.timestamp > cutoff]
|
|
cleared = len(self._records) - len(kept)
|
|
self._records = kept
|
|
self._error_counts = defaultdict(int)
|
|
self._plugin_error_counts = defaultdict(lambda: defaultdict(int))
|
|
for r in kept:
|
|
self._error_counts[r.error_type] += 1
|
|
if r.plugin_id:
|
|
self._plugin_error_counts[r.plugin_id][r.error_type] += 1
|
|
# A pattern that began after the cutoff is made only of kept
|
|
# errors; any other would carry cleared ones in its count.
|
|
self._patterns = {k: p for k, p in self._patterns.items() if p.first_seen > cutoff}
|
|
self._version += 1
|
|
return cleared
|
|
|
|
def build_snapshot(self) -> Dict[str, Any]:
|
|
"""A bounded, JSON-safe copy of the summary for another process.
|
|
|
|
The shape of get_error_summary() plus ``generated_at`` and
|
|
``plugin_health`` (get_plugin_health() for every plugin with errors).
|
|
Messages, stack traces and context are clipped so that one plugin
|
|
raising a huge exception cannot make the snapshot large.
|
|
"""
|
|
with self._lock:
|
|
summary = self.get_error_summary()
|
|
summary["generated_at"] = datetime.now().isoformat()
|
|
summary["recent_errors"] = [
|
|
_compact_record(r) for r in summary["recent_errors"][-_SNAPSHOT_RECENT_ERRORS:]
|
|
]
|
|
patterns = {}
|
|
for key, pattern in list(summary["active_patterns"].items())[:_SNAPSHOT_MAX_PATTERNS]:
|
|
pattern = dict(pattern)
|
|
pattern["affected_plugins"] = [
|
|
_clip(p, _SNAPSHOT_ID_CHARS) for p in pattern.get("affected_plugins", [])
|
|
][:_SNAPSHOT_MAX_AFFECTED_PLUGINS]
|
|
pattern["sample_messages"] = [
|
|
_redacted_clip(m, _SNAPSHOT_SAMPLE_CHARS) for m in pattern.get("sample_messages", [])
|
|
][:3]
|
|
patterns[key] = pattern
|
|
summary["active_patterns"] = patterns
|
|
health = {}
|
|
for plugin_id in summary["plugin_error_counts"]:
|
|
entry = self.get_plugin_health(plugin_id)
|
|
if entry["last_error"] is not None:
|
|
entry["last_error"] = _compact_record(entry["last_error"])
|
|
health[plugin_id] = entry
|
|
summary["plugin_health"] = health
|
|
# Round-trip so a context value the cache's encoder cannot handle is
|
|
# turned into a string here rather than failing the write.
|
|
return json.loads(json.dumps(summary, default=str))
|
|
|
|
def export_to_file(self, filepath: Path) -> None:
|
|
"""
|
|
Export error data to JSON file.
|
|
|
|
Args:
|
|
filepath: Path to export file
|
|
"""
|
|
with self._lock:
|
|
data = {
|
|
"exported_at": datetime.now().isoformat(),
|
|
"summary": self.get_error_summary(),
|
|
"all_records": [r.to_dict() for r in self._records]
|
|
}
|
|
filepath.parent.mkdir(parents=True, exist_ok=True)
|
|
filepath.write_text(json.dumps(data, indent=2))
|
|
self.logger.info(f"Exported error data to {filepath}")
|
|
|
|
def _auto_export(self) -> None:
|
|
"""Auto-export on pattern detection (if export_path configured)."""
|
|
if self.export_path:
|
|
try:
|
|
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
|
|
filepath = self.export_path / f"errors_{timestamp}.json"
|
|
self.export_to_file(filepath)
|
|
except Exception as e:
|
|
self.logger.error(f"Auto-export failed: {e}")
|
|
|
|
|
|
# Global singleton instance
|
|
_error_aggregator: Optional[ErrorAggregator] = None
|
|
_aggregator_lock = threading.Lock()
|
|
|
|
|
|
def get_error_aggregator(
|
|
max_records: int = 1000,
|
|
pattern_threshold: int = 5,
|
|
pattern_window_minutes: int = 60,
|
|
export_path: Optional[Path] = None
|
|
) -> ErrorAggregator:
|
|
"""
|
|
Get or create the global error aggregator instance.
|
|
|
|
Args:
|
|
max_records: Maximum records to keep (only used on first call)
|
|
pattern_threshold: Pattern detection threshold (only used on first call)
|
|
pattern_window_minutes: Pattern detection window (only used on first call)
|
|
export_path: Export path for auto-export (only used on first call)
|
|
|
|
Returns:
|
|
The global ErrorAggregator instance
|
|
"""
|
|
global _error_aggregator
|
|
|
|
with _aggregator_lock:
|
|
if _error_aggregator is None:
|
|
_error_aggregator = ErrorAggregator(
|
|
max_records=max_records,
|
|
pattern_threshold=pattern_threshold,
|
|
pattern_window_minutes=pattern_window_minutes,
|
|
export_path=export_path
|
|
)
|
|
return _error_aggregator
|
|
|
|
|
|
def record_error(
|
|
error: Exception,
|
|
context: Optional[Dict[str, Any]] = None,
|
|
plugin_id: Optional[str] = None,
|
|
operation: Optional[str] = None
|
|
) -> ErrorRecord:
|
|
"""
|
|
Convenience function to record an error to the global aggregator.
|
|
|
|
Args:
|
|
error: The exception that occurred
|
|
context: Optional context dictionary
|
|
plugin_id: Optional plugin ID
|
|
operation: Optional operation name
|
|
|
|
Returns:
|
|
The created ErrorRecord
|
|
"""
|
|
return get_error_aggregator().record_error(
|
|
error=error,
|
|
context=context,
|
|
plugin_id=plugin_id,
|
|
operation=operation
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Sharing the display service's errors with the web interface
|
|
# ---------------------------------------------------------------------------
|
|
#
|
|
# The two services are separate processes, so each has its own aggregator,
|
|
# and only the display service's ever records anything (plugin_executor runs
|
|
# the plugins there). The web interface therefore reads a snapshot the display
|
|
# service publishes to the shared cache directory -- the same channel, and the
|
|
# same file permissions, as display_current_state and plugin_metrics:*: files
|
|
# are 0660 and carry the cache directory's group, so root writes and the web
|
|
# user reads, and the other way round for the clear request.
|
|
#
|
|
# ERROR_SNAPSHOT_KEY written by the display service only
|
|
# ERROR_CLEAR_REQUEST_KEY written by the web interface only
|
|
#
|
|
# A clear is asynchronous: the web interface records a request, and the
|
|
# display service applies it (clear_before) on its next tick and republishes.
|
|
# Until it has, the web interface hides whatever the snapshot shows from
|
|
# before the cutoff, so a clear takes effect for readers immediately and a
|
|
# snapshot published just before the request cannot bring old errors back.
|
|
# The web interface never writes the snapshot itself: two writers would race,
|
|
# and a snapshot owned by the web user is one more file root's write has to
|
|
# replace.
|
|
|
|
ERROR_SNAPSHOT_KEY = "plugin_error_snapshot"
|
|
ERROR_CLEAR_REQUEST_KEY = "plugin_error_clear_request"
|
|
|
|
#: Shortest gap between two snapshot writes, in seconds. A plugin failing in
|
|
#: a tight loop changes the aggregator many times a second; the snapshot is
|
|
#: rewritten at most this often, and only when something changed.
|
|
SNAPSHOT_MIN_INTERVAL = 10.0
|
|
|
|
#: How often the display service checks for changes and clear requests. A
|
|
#: check is an in-memory comparison plus reading one small file.
|
|
SNAPSHOT_TICK_INTERVAL = 5.0
|
|
|
|
_SNAPSHOT_RECENT_ERRORS = 20
|
|
_SNAPSHOT_MAX_PATTERNS = 50
|
|
_SNAPSHOT_MAX_AFFECTED_PLUGINS = 50
|
|
_SNAPSHOT_MESSAGE_CHARS = 300
|
|
_SNAPSHOT_SAMPLE_CHARS = 200
|
|
_SNAPSHOT_TRACE_CHARS = 1200
|
|
_SNAPSHOT_ID_CHARS = 100
|
|
_SNAPSHOT_CONTEXT_KEYS = 20
|
|
_SNAPSHOT_CONTEXT_VALUE_CHARS = 200
|
|
|
|
#: Fields of get_error_summary(), which is what /errors/summary has always
|
|
#: returned. The snapshot's extra bookkeeping stays out of that response.
|
|
_SUMMARY_FIELDS = (
|
|
"session_start", "total_errors", "error_rate_per_hour",
|
|
"error_counts_by_type", "plugin_error_counts", "active_patterns",
|
|
"recent_errors",
|
|
)
|
|
|
|
_snapshot_logger = logging.getLogger(__name__ + ".snapshot")
|
|
|
|
|
|
def _clip(value: Any, limit: int) -> str:
|
|
text = value if isinstance(value, str) else str(value)
|
|
return text if len(text) <= limit else text[:limit - 3] + "..."
|
|
|
|
|
|
def _redacted_clip(value: Any, limit: int) -> str:
|
|
"""Redact, then clip. Clipping first could cut a ``token=`` marker off
|
|
while keeping the secret after it, and the web side's redaction would then
|
|
have nothing to match."""
|
|
return _clip(redact_credentials(value if isinstance(value, str) else str(value)), limit)
|
|
|
|
|
|
def _compact_record(record: Dict[str, Any]) -> Dict[str, Any]:
|
|
"""An ErrorRecord dict with every free-text field redacted and bounded."""
|
|
compact = dict(record)
|
|
compact["message"] = _redacted_clip(record.get("message") or "", _SNAPSHOT_MESSAGE_CHARS)
|
|
if record.get("plugin_id") is not None:
|
|
compact["plugin_id"] = _clip(record["plugin_id"], _SNAPSHOT_ID_CHARS)
|
|
trace = record.get("stack_trace")
|
|
if isinstance(trace, str):
|
|
trace = redact_credentials(trace)
|
|
if len(trace) > _SNAPSHOT_TRACE_CHARS:
|
|
# The end of a traceback is the part that says what went wrong.
|
|
trace = "..." + trace[-(_SNAPSHOT_TRACE_CHARS - 3):]
|
|
compact["stack_trace"] = trace
|
|
context = record.get("context")
|
|
if isinstance(context, dict):
|
|
compact["context"] = {
|
|
_clip(k, _SNAPSHOT_ID_CHARS): (
|
|
v if v is None or isinstance(v, (bool, int, float))
|
|
else _redacted_clip(v, _SNAPSHOT_CONTEXT_VALUE_CHARS)
|
|
)
|
|
for k, v in list(context.items())[:_SNAPSHOT_CONTEXT_KEYS]
|
|
}
|
|
else:
|
|
compact["context"] = {}
|
|
return compact
|
|
|
|
|
|
class ErrorSnapshotPublisher:
|
|
"""Publishes the display service's aggregator to the shared cache.
|
|
|
|
Runs in the display service only. tick() is the whole job; start() just
|
|
calls it from a daemon thread every SNAPSHOT_TICK_INTERVAL seconds, which
|
|
also means errors recorded while a write was being throttled still reach
|
|
the cache once the interval has passed, and a clear request is applied
|
|
even when no new error arrives to trigger a publish.
|
|
|
|
Nothing here raises: a failure to read or write the cache is logged at
|
|
debug and retried on a later tick.
|
|
"""
|
|
|
|
def __init__(self, cache_manager: Any, aggregator: Optional[ErrorAggregator] = None,
|
|
min_interval: float = SNAPSHOT_MIN_INTERVAL,
|
|
clock: Callable[[], float] = time.monotonic) -> None:
|
|
self.cache_manager = cache_manager
|
|
self.aggregator = aggregator or get_error_aggregator()
|
|
self.min_interval = min_interval
|
|
self._clock = clock
|
|
# None forces a first publish, which replaces a snapshot left behind
|
|
# by a previous run of the service with this run's (empty) one.
|
|
self._published_version: Optional[int] = None
|
|
self._last_attempt: Optional[float] = None
|
|
self._applied_clear_id: Optional[str] = None
|
|
self._tick_lock = threading.Lock()
|
|
self._stop = threading.Event()
|
|
self._thread: Optional[threading.Thread] = None
|
|
|
|
def _apply_clear_request(self) -> bool:
|
|
"""Honour a clear request we have not applied yet. True if one was."""
|
|
request = self.cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0)
|
|
if not isinstance(request, dict):
|
|
return False
|
|
request_id = request.get("request_id")
|
|
if not isinstance(request_id, str) or not request_id or request_id == self._applied_clear_id:
|
|
return False
|
|
try:
|
|
cutoff = float(request.get("cutoff"))
|
|
except (TypeError, ValueError):
|
|
cutoff = float("nan")
|
|
if math.isfinite(cutoff):
|
|
cleared = self.aggregator.clear_before(datetime.fromtimestamp(cutoff))
|
|
_snapshot_logger.info("Cleared %d plugin error record(s) as requested (%s)",
|
|
cleared, request_id)
|
|
# A malformed request is acknowledged too, so it is not retried forever.
|
|
self._applied_clear_id = request_id
|
|
return True
|
|
|
|
def tick(self) -> bool:
|
|
"""Apply a pending clear and publish if due. True if a snapshot was written."""
|
|
with self._tick_lock:
|
|
try:
|
|
cleared = self._apply_clear_request()
|
|
version = self.aggregator.version
|
|
now = self._clock()
|
|
if not cleared:
|
|
if version == self._published_version:
|
|
return False
|
|
if (self._last_attempt is not None
|
|
and now - self._last_attempt < self.min_interval):
|
|
return False
|
|
# Stamp the attempt before writing: a cache that keeps failing
|
|
# is retried at the throttled rate, not on every tick.
|
|
self._last_attempt = now
|
|
snapshot = self.aggregator.build_snapshot()
|
|
snapshot["applied_clear_id"] = self._applied_clear_id
|
|
self.cache_manager.set(ERROR_SNAPSHOT_KEY, snapshot)
|
|
self._published_version = version
|
|
return True
|
|
except Exception as err: # never let reporting break the display
|
|
_snapshot_logger.debug("Could not publish the plugin error snapshot: %s",
|
|
err, exc_info=True)
|
|
return False
|
|
|
|
def start(self, interval: float = SNAPSHOT_TICK_INTERVAL) -> None:
|
|
"""Tick from a daemon thread until stop(). A no-op while running."""
|
|
if self._thread is not None and self._thread.is_alive():
|
|
return
|
|
self._stop.clear()
|
|
|
|
def run() -> None:
|
|
self.tick()
|
|
while not self._stop.wait(interval):
|
|
self.tick()
|
|
|
|
self._thread = threading.Thread(target=run, name="error-snapshot-publisher", daemon=True)
|
|
self._thread.start()
|
|
|
|
def stop(self) -> None:
|
|
self._stop.set()
|
|
if self._thread is not None:
|
|
self._thread.join(timeout=2)
|
|
self._thread = None
|
|
|
|
|
|
_snapshot_publisher: Optional[ErrorSnapshotPublisher] = None
|
|
_snapshot_publisher_lock = threading.Lock()
|
|
|
|
|
|
def start_error_snapshot_publisher(cache_manager: Any) -> Optional[ErrorSnapshotPublisher]:
|
|
"""Start publishing this process's errors for the web interface.
|
|
|
|
Call from the display service only: whichever process calls it becomes
|
|
the source of /api/v3/errors/*. Idempotent; never raises.
|
|
"""
|
|
global _snapshot_publisher
|
|
try:
|
|
with _snapshot_publisher_lock:
|
|
if _snapshot_publisher is None:
|
|
_snapshot_publisher = ErrorSnapshotPublisher(cache_manager)
|
|
else:
|
|
_snapshot_publisher.cache_manager = cache_manager
|
|
_snapshot_publisher.start()
|
|
return _snapshot_publisher
|
|
except Exception as err:
|
|
_snapshot_logger.warning("Plugin error reporting to the web interface is unavailable: %s", err)
|
|
return None
|
|
|
|
|
|
# --- Reading side (web interface) -------------------------------------------
|
|
|
|
def read_error_report(cache_manager: Any) -> Tuple[Optional[Dict[str, Any]], Optional[Dict[str, Any]]]:
|
|
"""The display service's latest snapshot and the latest clear request.
|
|
|
|
memory_ttl=0: both keys are written by the other process, so only the
|
|
file is current.
|
|
"""
|
|
snapshot = cache_manager.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0)
|
|
clear_request = cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0)
|
|
return (snapshot if isinstance(snapshot, dict) else None,
|
|
clear_request if isinstance(clear_request, dict) else None)
|
|
|
|
|
|
def _epoch(iso: Any) -> Optional[float]:
|
|
"""Seconds since the epoch for an aggregator timestamp (local, naive)."""
|
|
if not isinstance(iso, str):
|
|
return None
|
|
try:
|
|
return datetime.fromisoformat(iso).timestamp()
|
|
except (ValueError, OverflowError, OSError):
|
|
return None
|
|
|
|
|
|
def _pending_cutoff(snapshot: Optional[Dict[str, Any]],
|
|
clear_request: Optional[Dict[str, Any]]) -> Optional[float]:
|
|
"""The cutoff of a clear the snapshot has not applied yet, if any."""
|
|
if not clear_request:
|
|
return None
|
|
request_id = clear_request.get("request_id")
|
|
if not request_id:
|
|
return None
|
|
if snapshot is not None and snapshot.get("applied_clear_id") == request_id:
|
|
return None
|
|
try:
|
|
cutoff = float(clear_request.get("cutoff"))
|
|
except (TypeError, ValueError):
|
|
return None
|
|
return cutoff if math.isfinite(cutoff) else None
|
|
|
|
|
|
def _is_after(item: Any, field_name: str, cutoff: float) -> bool:
|
|
when = _epoch(item.get(field_name)) if isinstance(item, dict) else None
|
|
return when is not None and when > cutoff
|
|
|
|
|
|
def _empty_summary(snapshot: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
|
return {
|
|
"session_start": snapshot.get("session_start") if snapshot else None,
|
|
"total_errors": 0,
|
|
"error_rate_per_hour": 0.0,
|
|
"error_counts_by_type": {},
|
|
"plugin_error_counts": {},
|
|
"active_patterns": {},
|
|
"recent_errors": [],
|
|
}
|
|
|
|
|
|
def error_summary_from_report(snapshot: Optional[Dict[str, Any]],
|
|
clear_request: Optional[Dict[str, Any]]) -> Dict[str, Any]:
|
|
"""The /errors/summary payload: get_error_summary()'s shape plus
|
|
``generated_at``, ``snapshot_available`` and ``clear_pending``."""
|
|
cutoff = _pending_cutoff(snapshot, clear_request)
|
|
summary = _empty_summary(snapshot)
|
|
if snapshot is not None:
|
|
for name, default in summary.items():
|
|
value = snapshot.get(name)
|
|
if isinstance(value, type(default)) or (
|
|
isinstance(default, float) and isinstance(value, int)) or (
|
|
name == "session_start" and isinstance(value, str)):
|
|
summary[name] = value
|
|
if cutoff is not None:
|
|
recent = summary["recent_errors"]
|
|
newest = _epoch(recent[-1].get("timestamp")) if recent and isinstance(recent[-1], dict) else None
|
|
if newest is None or newest <= cutoff:
|
|
# Everything the display has reported predates the clear.
|
|
summary = _empty_summary(snapshot)
|
|
else:
|
|
# Only part of it does. The lists can be filtered exactly; the
|
|
# counts cannot, and stay as reported until the display
|
|
# applies the clear (clear_pending says so).
|
|
summary["recent_errors"] = [r for r in recent if _is_after(r, "timestamp", cutoff)]
|
|
summary["active_patterns"] = {
|
|
k: p for k, p in summary["active_patterns"].items()
|
|
if _is_after(p, "last_seen", cutoff)
|
|
}
|
|
summary["generated_at"] = snapshot.get("generated_at") if snapshot else None
|
|
summary["snapshot_available"] = snapshot is not None
|
|
summary["clear_pending"] = cutoff is not None
|
|
return summary
|
|
|
|
|
|
def plugin_health_from_report(snapshot: Optional[Dict[str, Any]],
|
|
clear_request: Optional[Dict[str, Any]],
|
|
plugin_id: str) -> Dict[str, Any]:
|
|
"""The /errors/plugin/<id> payload: get_plugin_health()'s shape plus
|
|
``generated_at``, ``snapshot_available`` and ``clear_pending``."""
|
|
health: Dict[str, Any] = {
|
|
"plugin_id": plugin_id,
|
|
"status": "healthy",
|
|
"total_errors": 0,
|
|
"error_types": {},
|
|
"recent_error_count": 0,
|
|
"last_error": None,
|
|
}
|
|
cutoff = _pending_cutoff(snapshot, clear_request)
|
|
table = snapshot.get("plugin_health") if snapshot else None
|
|
entry = table.get(plugin_id) if isinstance(table, dict) else None
|
|
if isinstance(entry, dict):
|
|
# last_error is the plugin's newest error: if even that predates a
|
|
# pending clear, so does everything else the display reported for it.
|
|
if cutoff is None or _is_after(entry.get("last_error"), "timestamp", cutoff):
|
|
for name in ("status", "total_errors", "error_types", "recent_error_count", "last_error"):
|
|
if name in entry:
|
|
health[name] = entry[name]
|
|
health["generated_at"] = snapshot.get("generated_at") if snapshot else None
|
|
health["snapshot_available"] = snapshot is not None
|
|
health["clear_pending"] = cutoff is not None
|
|
return health
|
|
|
|
|
|
def _count_cleared(summary: Dict[str, Any], cutoff: float) -> Optional[int]:
|
|
"""How many of the reported errors a clear at ``cutoff`` hides, if known.
|
|
|
|
Exact when every reported error predates the cutoff (always the case for
|
|
a clear of everything) or when the report lists every error; otherwise
|
|
only the display service knows, and None says so.
|
|
"""
|
|
total = summary["total_errors"]
|
|
times = [_epoch(r.get("timestamp")) if isinstance(r, dict) else None
|
|
for r in summary["recent_errors"]]
|
|
if total == 0 or (times and times[-1] is not None and times[-1] <= cutoff):
|
|
return total
|
|
if len(times) >= total and None not in times:
|
|
return sum(1 for t in times if t <= cutoff)
|
|
return None
|
|
|
|
|
|
def request_error_clear(cache_manager: Any, cutoff: float) -> Dict[str, Any]:
|
|
"""Ask the display service to forget errors recorded at or before ``cutoff``.
|
|
|
|
Returns ``request_id``, ``cutoff`` (ISO, local time), ``cleared_count``
|
|
(see _count_cleared) and ``clear_requested``. Raises OSError when the
|
|
request did not reach the shared cache, since a cache without a usable
|
|
directory accepts set() and keeps nothing.
|
|
|
|
A request the display has not applied yet is only ever widened: a later,
|
|
narrower one ("older than 24 hours" after "everything") overwriting it
|
|
would otherwise bring back the errors the first one hid.
|
|
"""
|
|
snapshot, clear_request = read_error_report(cache_manager)
|
|
pending = _pending_cutoff(snapshot, clear_request)
|
|
if pending is not None:
|
|
cutoff = max(cutoff, pending)
|
|
before = error_summary_from_report(snapshot, clear_request)
|
|
request = {
|
|
"request_id": uuid.uuid4().hex,
|
|
"cutoff": cutoff,
|
|
"requested_at": time.time(),
|
|
}
|
|
cache_manager.set(ERROR_CLEAR_REQUEST_KEY, request)
|
|
stored = cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0)
|
|
if not isinstance(stored, dict) or stored.get("request_id") != request["request_id"]:
|
|
raise OSError("the clear request was not stored in the shared cache")
|
|
return {
|
|
"cleared_count": _count_cleared(before, cutoff),
|
|
"clear_requested": True,
|
|
"request_id": request["request_id"],
|
|
"cutoff": datetime.fromtimestamp(cutoff).isoformat(),
|
|
}
|