Files
LEDMatrix/src/error_aggregator.py
T
ChuckandClaude Opus 5.5 7b90759252 fix: /errors stack traces, Wi-Fi disconnect and save, plugin fonts, API cache TTL (#636)
* 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>
2026-09-24 17:32:29 -04:00

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(),
}