diff --git a/CHANGELOG.md b/CHANGELOG.md index 8b2d3927..72f9cbb9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -89,6 +89,23 @@ floor on the release that ships them): longer set `Accept-Encoding: ... br` by hand (brotli is not installed, so a `br` response could not be decoded); requests picks the encodings. +### Plugin error reporting + +- `/api/v3/errors/summary` and `/api/v3/errors/plugin/` report the errors + the display service recorded. They used to read the web process's own error + aggregator, which never records anything, so they always answered "no + errors". The display service now publishes a bounded snapshot to the shared + cache (`plugin_error_snapshot`, at most every 10 seconds and only on change; + `src/error_aggregator.py`, started from `DisplayController.__init__`). + Responses keep their shape and add `snapshot_available`, `generated_at` and + `clear_pending`; exception text has credentials redacted. +- `POST /api/v3/errors/clear` records a request (`plugin_error_clear_request`) + the display service applies within about 5 seconds; reads hide the cleared + errors at once. It accepts `"all": true`, and `cleared_count` can be `null` + when the count is only known to the display service. +- The Logs tab has a **Plugin errors** panel: per-plugin counts, repeating + errors and a Clear button. + ## 3.5.0 New modules a plugin may import via `src.*` (floor on 3.5.0): diff --git a/docs/PLUGIN_ERROR_HANDLING.md b/docs/PLUGIN_ERROR_HANDLING.md index a697dde8..16f02a09 100644 --- a/docs/PLUGIN_ERROR_HANDLING.md +++ b/docs/PLUGIN_ERROR_HANDLING.md @@ -146,7 +146,10 @@ def display(self, force_clear: bool = False) -> bool: ## Error Aggregation -LEDMatrix automatically tracks plugin errors. Access error data via the API: +LEDMatrix automatically tracks plugin errors: every exception or timeout from +a plugin's `update()` or `display()` is recorded by the display service, +which runs the plugins. See them in the web interface under **Logs → Plugin +errors**, or through the API: ```bash # Get error summary @@ -155,10 +158,21 @@ curl http://localhost:5000/api/v3/errors/summary # Get plugin-specific health curl http://localhost:5000/api/v3/errors/plugin/my-plugin -# Clear old errors +# Clear errors older than 24 hours (the default), or all of them curl -X POST http://localhost:5000/api/v3/errors/clear +curl -X POST -H 'Content-Type: application/json' -d '{"all": true}' \ + http://localhost:5000/api/v3/errors/clear ``` +The web interface is a separate process, so it reads a snapshot the display +service writes to the shared cache directory (`plugin_error_snapshot`): at most +every 10 seconds, and only when something changed. Expect the numbers to lag +by up to about 15 seconds, and to start from zero when the display service +restarts. `snapshot_available` is `false` until the display service has +reported. A clear is a request the display service applies within about 5 +seconds; the API hides the cleared errors immediately. Details and response +shapes: [REST API reference](REST_API_REFERENCE.md#error-tracking). + ### Error Patterns When the same error occurs repeatedly (5+ times in 60 minutes), it's detected as a pattern and logged as a warning. This helps identify systemic issues. diff --git a/docs/REST_API_REFERENCE.md b/docs/REST_API_REFERENCE.md index 26d2bb4d..f651a222 100644 --- a/docs/REST_API_REFERENCE.md +++ b/docs/REST_API_REFERENCE.md @@ -1815,25 +1815,73 @@ The last 100 journal lines for `ledmatrix.service` and ## Error tracking +Plugin errors are recorded by the display service (`ledmatrix.service`), +which runs the plugins. It publishes a snapshot to the shared cache directory +(`plugin_error_snapshot`) at most every 10 seconds, and only when something +changed, so these endpoints lag the display by up to about 15 seconds. The +counts cover the display service's current run: they start at zero when it +restarts. Error messages and stack traces have credentials redacted, and +messages, traces and context values are truncated in the snapshot. + +Every response below adds three fields to the shape it always had: + +| Field | Meaning | +|---|---| +| `snapshot_available` | `false` until the display service has reported (for example, it is not running). Counts are then zero. | +| `generated_at` | When the display service produced the snapshot (ISO, the Pi's local time), or `null`. | +| `clear_pending` | A clear has been requested and the display service has not applied it yet. | + ### Get Error Summary **GET** `/api/v3/errors/summary` -Aggregated counts, detected patterns and recent errors across plugins and -core components. +Aggregated counts, detected patterns and recent errors (the last 20). + +```json +{ + "status": "success", + "data": { + "session_start": "2026-09-23T09:40:02.118000", + "total_errors": 13, + "error_rate_per_hour": 41.2, + "error_counts_by_type": {"ConnectionError": 12, "ValueError": 1}, + "plugin_error_counts": {"weather": {"ConnectionError": 12}, "stocks": {"ValueError": 1}}, + "active_patterns": { + "ConnectionError": { + "error_type": "ConnectionError", "count": 12, + "first_seen": "2026-09-23T09:41:10.500000", "last_seen": "2026-09-23T09:58:36.020000", + "affected_plugins": ["weather"], "sample_messages": ["Read timed out."], + "severity": "error" + } + }, + "recent_errors": [ + {"error_type": "ValueError", "message": "could not parse price", + "timestamp": "2026-09-23T09:58:36.100000", "context": {}, + "plugin_id": "stocks", "operation": "update", "stack_trace": "Traceback ..."} + ], + "generated_at": "2026-09-23T09:58:40.000000", + "snapshot_available": true, + "clear_pending": false + }, + "message": "Error summary retrieved" +} +``` ### Get Plugin Errors **GET** `/api/v3/errors/plugin/` -Error health and statistics for one plugin. +Error health and statistics for one plugin: `plugin_id`, `status` +(`healthy`, `degraded` or `unhealthy`), `total_errors`, `error_types`, +`recent_error_count`, `last_error` (a `recent_errors` entry or `null`), plus +the three fields above. A plugin with no recorded errors is `healthy`. ### Clear Errors **POST** `/api/v3/errors/clear` -Clear error records older than `max_age_hours` (default 24, 1-8760). -Returns `data.cleared_count`. +Clear error records older than `max_age_hours` (default 24, 1-8760), or every +error with `"all": true` (`max_age_hours` is then ignored). ```json { @@ -1841,6 +1889,32 @@ Returns `data.cleared_count`. } ``` +The clear is asynchronous. The web interface records a request +(`plugin_error_clear_request` in the shared cache), and the display service +applies it within about 5 seconds, rebuilding its counts from the errors it +keeps and republishing. Reads hide the cleared errors from the moment the +request is recorded. Until the display service applies an age-based clear, +`recent_errors` and `active_patterns` are already filtered but the counts +are the old ones, and `clear_pending` is `true`. + +```json +{ + "status": "success", + "data": { + "cleared_count": 13, + "clear_requested": true, + "request_id": "5f0c1e...", + "cutoff": "2026-09-23T09:59:02.310000" + }, + "message": "Clear of all errors requested; the display service applies it within about 5 seconds" +} +``` + +`cleared_count` is how many of the reported errors the clear hides. It is +`null` when that cannot be known before the display service applies it (an +age-based clear over more errors than the report lists). A request that +could not be written to the shared cache answers `500`. + --- ## Health and Status diff --git a/src/display_controller.py b/src/display_controller.py index 1124f01b..1385cb82 100644 --- a/src/display_controller.py +++ b/src/display_controller.py @@ -95,6 +95,9 @@ class DisplayController: self.config_manager = config_manager # Keep for backward compatibility self.config = self.config_service.get_config() self.cache_manager = CacheManager() + # The web interface's /api/v3/errors/* read what this publishes. + from src.error_aggregator import start_error_snapshot_publisher + start_error_snapshot_publisher(self.cache_manager) logger.info("Config loaded in %.3f seconds (hot-reload: %s)", time.time() - start_time, enable_hot_reload) # Validate startup configuration diff --git a/src/error_aggregator.py b/src/error_aggregator.py index 42eb65c7..55fc7b1e 100644 --- a/src/error_aggregator.py +++ b/src/error_aggregator.py @@ -9,17 +9,21 @@ 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 +from typing import Dict, List, Optional, Any, Callable, Tuple import logging from src.exceptions import LEDMatrixError +from src.redaction import redact_credentials @dataclass @@ -115,6 +119,10 @@ class ErrorAggregator: # 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, @@ -161,6 +169,7 @@ class ErrorAggregator: 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) @@ -331,6 +340,77 @@ class ErrorAggregator: 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. @@ -419,3 +499,391 @@ def record_error( 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/ 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(), + } diff --git a/src/redaction.py b/src/redaction.py new file mode 100644 index 00000000..b84bf1b1 --- /dev/null +++ b/src/redaction.py @@ -0,0 +1,49 @@ +"""Credential redaction for text that leaves the process that produced it. + +Kept free of Flask so the display service can redact what it publishes (see +src/error_aggregator.py) as well as the web interface what it returns. +""" + +import re + +# Credentials that turn up inside exception text. A requests error quotes the +# URL it failed on, and plugins that authenticate by query string put their key +# there, so echoing an exception verbatim can hand out an API key. Redact the +# value, keep the parameter name -- knowing *which* credential was involved is +# part of the diagnosis. +_REDACT_CREDENTIAL = re.compile( + r'((?:api[_-]?key|access[_-]?token|auth|apikey|key|passwd|password|pwd|' + r'secret|sig|signature|token)["\']?\s*[=:]\s*["\']?)([^\s&"\'<>,}]+)', + re.IGNORECASE, +) + +# `Authorization: `. The scheme name is kept because it +# says which kind of credential failed; the credential goes. Any scheme +# matches, not a fixed list: ApiKey, Negotiate, NTLM, AWS4-HMAC-SHA256 and +# whatever a plugin's API invents next are all credentials, and a list would +# silently leak the ones nobody thought of. Not covered by the generic pattern +# above, whose value part stops at whitespace and so would keep the credential +# once a space follows the scheme. +_REDACT_AUTH_HEADER = re.compile( + r'((?:proxy-)?authorization["\']?\s*[=:]\s*["\']?\s*' + r'(?:[A-Za-z][\w.+-]*[ \t]+)?)' # optional scheme name, kept + r'([^\s,"\'<>}]+)', # the credential, redacted + re.IGNORECASE, +) + +# Credentials embedded in a URL: https://user:password@host. requests quotes +# the full URL in its exceptions, so this is a realistic leak. The username is +# kept -- it identifies which account failed without being the secret. +_REDACT_URL_USERINFO = re.compile(r'([a-z][a-z0-9+.-]*://[^/\s:@]+:)([^/\s@]+)(@)', + re.IGNORECASE) + + +def redact_credentials(text: str) -> str: + """Replace credentials in ``text`` with ````; keep everything + else, including line breaks, so a stack trace stays readable.""" + text = text or '' + # Order matters: the URL and header forms are more specific than the + # generic key=value pattern, which would otherwise chew the scheme. + text = _REDACT_URL_USERINFO.sub(r'\1\3', text) + text = _REDACT_AUTH_HEADER.sub(r'\1', text) + return _REDACT_CREDENTIAL.sub(r'\1', text) diff --git a/src/web_interface/error_handler.py b/src/web_interface/error_handler.py index 509c5679..6554568c 100644 --- a/src/web_interface/error_handler.py +++ b/src/web_interface/error_handler.py @@ -4,7 +4,6 @@ Centralized error handling for web interface. Provides helpers for consistent error responses across API endpoints. """ -import re from typing import Any, Optional from flask import jsonify @@ -12,42 +11,12 @@ from src.web_interface.errors import ( WebInterfaceError, ErrorCode, ErrorCategory ) from src.logging_config import get_logger +from src.redaction import redact_credentials logger = get_logger(__name__) -# Credentials that turn up inside exception text. A requests error quotes the -# URL it failed on, and plugins that authenticate by query string put their key -# there, so echoing an exception verbatim can hand out an API key. Redact the -# value, keep the parameter name -- knowing *which* credential was involved is -# part of the diagnosis. -_REDACT_CREDENTIAL = re.compile( - r'((?:api[_-]?key|access[_-]?token|auth|apikey|key|passwd|password|pwd|' - r'secret|sig|signature|token)["\']?\s*[=:]\s*["\']?)([^\s&"\'<>,}]+)', - re.IGNORECASE, -) - -# `Authorization: `. The scheme name is kept because it -# says which kind of credential failed; the credential goes. Any scheme -# matches, not a fixed list: ApiKey, Negotiate, NTLM, AWS4-HMAC-SHA256 and -# whatever a plugin's API invents next are all credentials, and a list would -# silently leak the ones nobody thought of. Not covered by the generic pattern -# above, whose value part stops at whitespace and so would keep the credential -# once a space follows the scheme. -_REDACT_AUTH_HEADER = re.compile( - r'((?:proxy-)?authorization["\']?\s*[=:]\s*["\']?\s*' - r'(?:[A-Za-z][\w.+-]*[ \t]+)?)' # optional scheme name, kept - r'([^\s,"\'<>}]+)', # the credential, redacted - re.IGNORECASE, -) - -# Credentials embedded in a URL: https://user:password@host. requests quotes -# the full URL in its exceptions, so this is a realistic leak. The username is -# kept -- it identifies which account failed without being the secret. -_REDACT_URL_USERINFO = re.compile(r'([a-z][a-z0-9+.-]*://[^/\s:@]+:)([^/\s@]+)(@)', - re.IGNORECASE) - # Long enough for an errno string with a path, short enough not to dump a # parser's worth of context into a JSON field. _MAX_DETAIL_LENGTH = 400 @@ -95,12 +64,7 @@ def redact_text(text: str, max_length: int = _MAX_DETAIL_LENGTH) -> str: Returns: A single line, credentials replaced, length capped. """ - text = text or '' - # Order matters: the URL and header forms are more specific than the - # generic key=value pattern, which would otherwise chew the scheme. - text = _REDACT_URL_USERINFO.sub(r'\1\3', text) - text = _REDACT_AUTH_HEADER.sub(r'\1', text) - text = _REDACT_CREDENTIAL.sub(r'\1', text) + text = redact_credentials(text) # Collapse newlines/tabs so the detail stays one line in a JSON field. text = ' '.join(text.split()) if len(text) > max_length: diff --git a/test/test_error_snapshot_cross_process.py b/test/test_error_snapshot_cross_process.py new file mode 100644 index 00000000..3cc48614 --- /dev/null +++ b/test/test_error_snapshot_cross_process.py @@ -0,0 +1,443 @@ +"""/api/v3/errors/* report the display service's errors, not the web's. + +The error aggregator is a per-process singleton, and only the display service +runs plugins, so only its aggregator ever records anything. The routes used to +read the web process's own aggregator and so always answered "no errors". + +Here the two services are two ErrorAggregator instances and two CacheManagers +over one temporary directory -- the same arrangement as the real services, +which share /var/cache/ledmatrix. The display side publishes through +ErrorSnapshotPublisher.tick(); the web side is the real blueprint. +""" +import json +import os +import stat +import sys +from datetime import datetime, timedelta +from pathlib import Path +from unittest.mock import MagicMock + +import pytest + +sys.path.insert(0, str(Path(__file__).parent.parent)) + +from src.cache_manager import CacheManager # noqa: E402 +from src import error_aggregator as errors # noqa: E402 +from src.error_aggregator import ( # noqa: E402 + ERROR_CLEAR_REQUEST_KEY, ERROR_SNAPSHOT_KEY, ErrorAggregator, + ErrorSnapshotPublisher, +) +from test._api_v3_test_helpers import api_v3_client, api_v3_module # noqa: F401,E402 + + +class FakeClock: + def __init__(self): + self.now = 1000.0 + + def __call__(self): + return self.now + + +def _fail(aggregator, plugin_id="p1", message="boom", exc=ValueError, **kwargs): + try: + raise exc(message) + except Exception as e: # noqa: BLE001 - recording is the point + return aggregator.record_error(e, plugin_id=plugin_id, operation="update", **kwargs) + + +@pytest.fixture +def shared_cache(tmp_path, monkeypatch): + """Two cache managers over one directory: the display's and the web's.""" + monkeypatch.setattr(CacheManager, "_get_writable_cache_dir", lambda self: str(tmp_path)) + display_cache, web_cache = CacheManager(), CacheManager() + yield display_cache, web_cache, tmp_path + display_cache.stop_cleanup_thread() + web_cache.stop_cleanup_thread() + + +@pytest.fixture +def display(shared_cache): + display_cache, _, _ = shared_cache + aggregator = ErrorAggregator() + clock = FakeClock() + publisher = ErrorSnapshotPublisher(display_cache, aggregator=aggregator, clock=clock) + return aggregator, publisher, clock + + +@pytest.fixture +def web(api_v3_module, api_v3_client, shared_cache): # noqa: F811 + _, web_cache, _ = shared_cache + api_v3_module.api_v3.cache_manager = web_cache + return api_v3_client + + +def _summary(client): + response = client.get("/api/v3/errors/summary") + assert response.status_code == 200, response.get_data(as_text=True) + return response.get_json()["data"] + + +# --- Display side ----------------------------------------------------------- + +class TestPublishing: + def test_errors_reach_the_shared_cache(self, display, shared_cache): + aggregator, publisher, _ = display + _, web_cache, _ = shared_cache + _fail(aggregator) + assert publisher.tick() is True + snapshot = web_cache.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0) + assert snapshot["total_errors"] == 1 + assert snapshot["plugin_error_counts"] == {"p1": {"ValueError": 1}} + assert datetime.fromisoformat(snapshot["generated_at"]) + + def test_first_tick_publishes_even_with_no_errors(self, display, shared_cache): + # Replaces whatever a previous run of the display service left behind. + _, web_cache, _ = shared_cache + web_cache.set(ERROR_SNAPSHOT_KEY, {"total_errors": 99}) + _, publisher, _ = display + assert publisher.tick() is True + assert web_cache.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0)["total_errors"] == 0 + + def test_nothing_is_written_when_nothing_changed(self, display, shared_cache): + aggregator, publisher, clock = display + publisher.tick() + publisher.cache_manager = MagicMock(wraps=publisher.cache_manager) + clock.now += 3600 + assert publisher.tick() is False + publisher.cache_manager.set.assert_not_called() + + def test_a_tight_failure_loop_is_throttled(self, display): + aggregator, publisher, clock = display + _fail(aggregator) + assert publisher.tick() is True + publisher.cache_manager = MagicMock(wraps=publisher.cache_manager) + for _ in range(50): + _fail(aggregator) + clock.now += 0.1 + publisher.tick() + publisher.cache_manager.set.assert_not_called() + # Once the interval has passed, the backlog is published without a + # new error having to arrive. + clock.now += errors.SNAPSHOT_MIN_INTERVAL + assert publisher.tick() is True + assert publisher.cache_manager.set.call_count == 1 + assert publisher.cache_manager.set.call_args[0][1]["total_errors"] == 51 + + def test_a_failing_cache_never_raises(self, display): + aggregator, publisher, clock = display + broken = MagicMock() + broken.get.side_effect = OSError("read-only file system") + broken.set.side_effect = OSError("read-only file system") + publisher.cache_manager = broken + _fail(aggregator) + assert publisher.tick() is False + broken.get.side_effect = None + broken.get.return_value = None + assert publisher.tick() is False # set still fails + # ...and retries at the throttled rate, not every tick. + calls = broken.set.call_count + clock.now += 1 + publisher.tick() + assert broken.set.call_count == calls + + def test_an_unserialisable_context_does_not_break_the_snapshot(self, display, shared_cache): + aggregator, publisher, _ = display + _fail(aggregator, context={"obj": object(), "n": 3}) + assert publisher.tick() is True + _, web_cache, _ = shared_cache + snap = web_cache.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0) + assert snap["recent_errors"][0]["context"]["n"] == 3 + + def test_snapshot_stays_small(self, display): + aggregator, _, _ = display + for i in range(300): + _fail(aggregator, plugin_id=f"plugin-{i % 5}", message="x" * 20000, + context={f"k{j}": "v" * 5000 for j in range(50)}) + snapshot = aggregator.build_snapshot() + assert len(snapshot["recent_errors"]) == 20 + assert all(len(r["message"]) <= 300 for r in snapshot["recent_errors"]) + assert len(json.dumps(snapshot)) < 200_000 + + def test_start_publishes_from_a_background_thread(self, shared_cache): + display_cache, web_cache, _ = shared_cache + aggregator = ErrorAggregator() + _fail(aggregator) + publisher = ErrorSnapshotPublisher(display_cache, aggregator=aggregator) + try: + publisher.start(interval=0.05) + deadline = datetime.now() + timedelta(seconds=5) + snapshot = None + while snapshot is None and datetime.now() < deadline: + snapshot = web_cache.get(ERROR_SNAPSHOT_KEY, max_age=None, memory_ttl=0) + assert snapshot and snapshot["total_errors"] == 1 + finally: + publisher.stop() + + def test_start_helper_never_raises(self, monkeypatch): + monkeypatch.setattr(errors, "_snapshot_publisher", None) + + def explode(*a, **k): + raise RuntimeError("no threads for you") + + monkeypatch.setattr(errors.ErrorSnapshotPublisher, "start", explode) + assert errors.start_error_snapshot_publisher(MagicMock()) is None + + def test_display_controller_starts_the_publisher(self): + # The display service is the only process that runs plugins, so it is + # the one that must publish; the web process must not. + # Read, not imported: display_controller needs rgbmatrix. + root = Path(__file__).parent.parent + controller = (root / "src" / "display_controller.py").read_text(encoding="utf-8") + init = controller.split(" def __init__(self):", 1)[1].split("\n def ", 1)[0] + assert "start_error_snapshot_publisher(self.cache_manager)" in init + web_dir = root / "web_interface" + for source in web_dir.rglob("*.py"): + assert "start_error_snapshot_publisher" not in source.read_text(encoding="utf-8"), source + + +class TestClearBefore: + def test_keeps_later_errors_and_rebuilds_counts(self): + aggregator = ErrorAggregator(pattern_threshold=2) + for _ in range(3): + _fail(aggregator, plugin_id="old") + for record in aggregator._records: + record.timestamp -= timedelta(hours=2) + for pattern in aggregator._patterns.values(): + pattern.first_seen -= timedelta(hours=2) + _fail(aggregator, plugin_id="new", exc=KeyError) + cleared = aggregator.clear_before(datetime.now() - timedelta(hours=1)) + assert cleared == 3 + summary = aggregator.get_error_summary() + assert summary["total_errors"] == 1 + assert summary["error_counts_by_type"] == {"KeyError": 1} + assert summary["plugin_error_counts"] == {"new": {"KeyError": 1}} + assert summary["active_patterns"] == {} + + def test_a_pattern_made_only_of_later_errors_survives(self): + # A clear applied a few seconds late must not drop a pattern that + # formed entirely after the cutoff. + aggregator = ErrorAggregator(pattern_threshold=2) + cutoff = datetime.now() - timedelta(seconds=1) + for _ in range(3): + _fail(aggregator) + aggregator.clear_before(cutoff) + assert list(aggregator.get_error_summary()["active_patterns"]) == ["ValueError"] + + def test_changes_the_version(self): + aggregator = ErrorAggregator() + v = aggregator.version + aggregator.clear_before(datetime.now()) + assert aggregator.version != v + + +# --- Web side --------------------------------------------------------------- + +class TestRoutes: + def test_no_snapshot_yet(self, web): + data = _summary(web) + assert data["snapshot_available"] is False + assert data["generated_at"] is None + assert data["total_errors"] == 0 + assert data["recent_errors"] == [] and data["active_patterns"] == {} + response = web.get("/api/v3/errors/summary") + assert "not reported" in response.get_json()["message"] + + def test_summary_is_the_display_services(self, web, display): + aggregator, publisher, _ = display + for _ in range(5): + _fail(aggregator, plugin_id="weather", message="HTTP 500") + publisher.tick() + data = _summary(web) + assert data["snapshot_available"] is True + assert data["clear_pending"] is False + assert data["total_errors"] == 5 + assert data["plugin_error_counts"] == {"weather": {"ValueError": 5}} + pattern = data["active_patterns"]["ValueError"] + assert pattern["affected_plugins"] == ["weather"] + assert pattern["sample_messages"] == ["HTTP 500"] + # The documented shape, plus only the documented additions. + assert set(data) == set(ErrorAggregator().get_error_summary()) | { + "generated_at", "snapshot_available", "clear_pending"} + + def test_web_process_aggregator_is_not_what_is_reported(self, web, display, monkeypatch): + # The original bug: the route read this process's own aggregator. + local = ErrorAggregator() + _fail(local, plugin_id="only-in-web") + monkeypatch.setattr(errors, "_error_aggregator", local) + aggregator, publisher, _ = display + _fail(aggregator, plugin_id="from-display") + publisher.tick() + assert list(_summary(web)["plugin_error_counts"]) == ["from-display"] + + def test_plugin_slice(self, web, display): + aggregator, publisher, _ = display + for _ in range(6): + _fail(aggregator, plugin_id="weather") + _fail(aggregator, plugin_id="clock", exc=KeyError) + publisher.tick() + data = web.get("/api/v3/errors/plugin/weather").get_json()["data"] + assert data["plugin_id"] == "weather" + assert data["status"] == "unhealthy" + assert data["total_errors"] == 6 + assert data["error_types"] == {"ValueError": 6} + assert data["last_error"]["plugin_id"] == "weather" + assert data["snapshot_available"] is True + healthy = web.get("/api/v3/errors/plugin/never-failed").get_json()["data"] + assert healthy["status"] == "healthy" and healthy["total_errors"] == 0 + assert healthy["last_error"] is None + + def test_credentials_in_exception_text_are_redacted(self, web, display): + aggregator, publisher, _ = display + for _ in range(5): + _fail(aggregator, message="GET https://api.example.com/?api_key=SEKRIT123 failed") + publisher.tick() + body = web.get("/api/v3/errors/summary").get_data(as_text=True) + body += web.get("/api/v3/errors/plugin/p1").get_data(as_text=True) + assert "SEKRIT123" not in body + assert "api_key=" in body + + def test_snapshot_is_redacted_before_it_is_clipped(self, display): + """The display redacts what it publishes, before clipping: keeping + only a traceback's tail could cut ``api_key=`` off and leave the key + itself, which the web side's redaction would then not recognise.""" + aggregator, _, _ = display + record = _fail(aggregator, message="GET /?token=MSGSECRET failed") + tail = errors._SNAPSHOT_TRACE_CHARS - 3 + filler = "x" * (tail - len("TRACESECRET ")) + record.stack_trace = "requests failed: api_key=TRACESECRET " + filler + assert record.stack_trace[-tail:].startswith("TRACESECRET") # marker falls outside + record.context = {"url": "https://h/?password=CTXSECRET"} + for _ in range(5): + _fail(aggregator, message="GET /?token=SAMPLESECRET failed") + published = json.dumps(aggregator.build_snapshot()) + for secret in ("MSGSECRET", "TRACESECRET", "CTXSECRET", "SAMPLESECRET"): + assert secret not in published, secret + + +class TestClear: + def test_clear_is_applied_by_the_display_and_republished(self, web, display): + aggregator, publisher, _ = display + _fail(aggregator) + _fail(aggregator) + publisher.tick() + response = web.post("/api/v3/errors/clear", json={"all": True}) + assert response.status_code == 200 + body = response.get_json()["data"] + assert body["clear_requested"] is True + assert body["cleared_count"] == 2 + # The display applies it on its next tick, throttle or not... + assert publisher.tick() is True + assert aggregator.get_error_summary()["total_errors"] == 0 + data = _summary(web) + assert data["total_errors"] == 0 and data["clear_pending"] is False + # ...and only once. + assert publisher.tick() is False + + def test_summary_hides_cleared_errors_before_the_display_applies_it(self, web, display): + aggregator, publisher, _ = display + for _ in range(5): + _fail(aggregator) + publisher.tick() + web.post("/api/v3/errors/clear", json={"all": True}) + # No display tick yet. + data = _summary(web) + assert data["clear_pending"] is True + assert data["total_errors"] == 0 + assert data["recent_errors"] == [] and data["active_patterns"] == {} + assert data["plugin_error_counts"] == {} + plugin = web.get("/api/v3/errors/plugin/p1").get_json()["data"] + assert plugin["status"] == "healthy" and plugin["total_errors"] == 0 + + def test_a_snapshot_written_just_before_the_clear_cannot_bring_errors_back( + self, web, display, shared_cache): + # The race: the display builds a snapshot, the user clicks Clear, and + # the display's write lands after the request. + aggregator, publisher, _ = display + display_cache, _, _ = shared_cache + for _ in range(3): + _fail(aggregator) + stale = aggregator.build_snapshot() + web.post("/api/v3/errors/clear", json={"all": True}) + stale["applied_clear_id"] = None + display_cache.set(ERROR_SNAPSHOT_KEY, stale) + assert _summary(web)["total_errors"] == 0 + + def test_errors_after_the_clear_are_kept(self, web, display): + aggregator, publisher, clock = display + _fail(aggregator, plugin_id="before") + publisher.tick() + web.post("/api/v3/errors/clear", json={"all": True}) + # An error lands after the request but before the display applies it. + for record in aggregator._records: + record.timestamp -= timedelta(seconds=5) + _fail(aggregator, plugin_id="after") + publisher.tick() + data = _summary(web) + assert data["plugin_error_counts"] == {"after": {"ValueError": 1}} + assert data["clear_pending"] is False + + def test_age_based_clear(self, web, display): + aggregator, publisher, _ = display + _fail(aggregator, plugin_id="old") + _fail(aggregator, plugin_id="old") + for record in aggregator._records: + record.timestamp -= timedelta(hours=3) + _fail(aggregator, plugin_id="new") + publisher.tick() + body = web.post("/api/v3/errors/clear", json={"max_age_hours": 1}).get_json()["data"] + assert body["cleared_count"] == 2 + pending = _summary(web) + assert pending["clear_pending"] is True + assert [r["plugin_id"] for r in pending["recent_errors"]] == ["new"] + publisher.tick() + data = _summary(web) + assert data["plugin_error_counts"] == {"new": {"ValueError": 1}} + assert data["total_errors"] == 1 + + def test_a_narrower_clear_does_not_undo_a_pending_wider_one(self, web, display): + aggregator, publisher, _ = display + for _ in range(3): + _fail(aggregator) + publisher.tick() + web.post("/api/v3/errors/clear", json={"all": True}) + web.post("/api/v3/errors/clear", json={"max_age_hours": 24}) + assert _summary(web)["total_errors"] == 0 + publisher.tick() + assert aggregator.get_error_summary()["total_errors"] == 0 + + def test_default_body_still_means_older_than_24_hours(self, web, display): + aggregator, publisher, _ = display + _fail(aggregator) + publisher.tick() + response = web.post("/api/v3/errors/clear") + assert response.status_code == 200 + assert response.get_json()["data"]["cleared_count"] == 0 + publisher.tick() + assert _summary(web)["total_errors"] == 1 + + def test_validation_is_unchanged_but_all_skips_it(self, web): + assert web.post("/api/v3/errors/clear", json={"max_age_hours": 0}).status_code == 400 + assert web.post("/api/v3/errors/clear", json={"max_age_hours": 9000}).status_code == 400 + assert web.post("/api/v3/errors/clear", + json={"all": True, "max_age_hours": "junk"}).status_code == 200 + + def test_a_request_that_did_not_reach_the_cache_is_an_error(self, web, api_v3_module): # noqa: F811 + cache = MagicMock() + cache.get.return_value = None + api_v3_module.api_v3.cache_manager = cache + response = web.post("/api/v3/errors/clear", json={"all": True}) + assert response.status_code == 500 + assert "clear request" in response.get_json()["message"] + + +@pytest.mark.skipif(not hasattr(os, "fchmod") or os.name == "nt", + reason="POSIX file modes") +def test_both_files_are_group_readable(web, display, shared_cache): + aggregator, publisher, _ = display + _, _, directory = shared_cache + _fail(aggregator) + publisher.tick() + web.post("/api/v3/errors/clear", json={"all": True}) + for key in (ERROR_SNAPSHOT_KEY, ERROR_CLEAR_REQUEST_KEY): + mode = stat.S_IMODE(os.stat(directory / f"{key}.json").st_mode) + assert mode == 0o660, (key, oct(mode)) diff --git a/web_interface/blueprints/api_v3/misc.py b/web_interface/blueprints/api_v3/misc.py index baeddf59..ddcb8790 100644 --- a/web_interface/blueprints/api_v3/misc.py +++ b/web_interface/blueprints/api_v3/misc.py @@ -10,10 +10,11 @@ from web_interface.blueprints.api_v3 import ( _MQTT_BRIDGE_DIR, _SUDO, _coerce_mqtt_bridge_value, _get_display_service_status, _mqtt_bridge_service_state, _read_mqtt_bridge_config, api_v3, contextlib, describe_exception, - error_response, get_error_aggregator, json, jsonify, logger, os, request, + error_response, json, jsonify, logger, os, redact_text, request, subprocess, success_response, tempfile, ) from src.common.path_safety import safe_path_component +from src import error_aggregator as _errors import web_interface.blueprints.api_v3 as _pkg # Read through the module rather than bound by value: tests patch these # as module attributes, and a value binding would not see the patch. @@ -326,17 +327,62 @@ def delete_cache_file(): except Exception as e: logger.error('Error in delete_cache_file', exc_info=True) return jsonify({'status': 'error', 'message': 'An error occurred; see logs for details', 'details': describe_exception(e)}), 500 +def _errors_cache(): + """The shared cache the display service publishes its errors to.""" + if not api_v3.cache_manager: + from src.cache_manager import CacheManager + api_v3.cache_manager = CacheManager() + return api_v3.cache_manager + + +def _redact_error_text(text, keep_lines=False): + """Credentials out of plugin exception text, which can quote a URL with + an API key in it. Stack traces keep their line breaks and indentation.""" + if not isinstance(text, str): + return text + if not keep_lines: + return redact_text(text, max_length=len(text) + 1) + return '\n'.join( + line[:len(line) - len(line.lstrip())] + redact_text(line, max_length=len(line) + 1) + for line in text.splitlines() + ) + + +def _redact_error_record(record): + if not isinstance(record, dict): + return record + record = dict(record) + record['message'] = _redact_error_text(record.get('message')) + record['stack_trace'] = _redact_error_text(record.get('stack_trace'), keep_lines=True) + if isinstance(record.get('context'), dict): + record['context'] = {k: _redact_error_text(v) for k, v in record['context'].items()} + return record + + +def _read_errors(): + snapshot, clear_request = _errors.read_error_report(_errors_cache()) + return snapshot, clear_request + + @api_v3.route('/errors/summary', methods=['GET']) def get_error_summary(): """ Get summary of all errors for monitoring and debugging. - Returns error counts, detected patterns, and recent errors. + Returns error counts, detected patterns, and recent errors, as last + reported by the display service (which runs the plugins, so it is the + only process that records their errors). ``snapshot_available`` is false + until it has reported; ``generated_at`` says when it did. """ try: - aggregator = get_error_aggregator() - summary = aggregator.get_error_summary() - return success_response(data=summary, message="Error summary retrieved") + summary = _errors.error_summary_from_report(*_read_errors()) + summary['recent_errors'] = [_redact_error_record(r) for r in summary['recent_errors']] + for pattern in summary['active_patterns'].values(): + if isinstance(pattern, dict) and isinstance(pattern.get('sample_messages'), list): + pattern['sample_messages'] = [_redact_error_text(m) for m in pattern['sample_messages']] + message = ("Error summary retrieved" if summary['snapshot_available'] + else "The display service has not reported any errors yet") + return success_response(data=summary, message=message) except Exception as e: logger.error(f"Error getting error summary: {e}", exc_info=True) return error_response( @@ -352,11 +398,13 @@ def get_plugin_errors(plugin_id): Args: plugin_id: Plugin identifier - Returns health status and error statistics for the plugin. + Returns health status and error statistics for the plugin, from the + display service's last report (see get_error_summary). A plugin with no + recorded errors is "healthy". """ try: - aggregator = get_error_aggregator() - health = aggregator.get_plugin_health(plugin_id) + health = _errors.plugin_health_from_report(*_read_errors(), plugin_id) + health['last_error'] = _redact_error_record(health['last_error']) return success_response(data=health, message="Plugin health retrieved") except Exception as e: logger.error(f"Error getting plugin health for {plugin_id}: {e}", exc_info=True) @@ -372,42 +420,61 @@ def clear_old_errors(): Request body (optional): max_age_hours: Maximum age in hours (default: 24, max: 8760 = 1 year) + all: true clears every error recorded so far (max_age_hours ignored) + + The errors live in the display service, so this records a clear request + that it applies within a few seconds. Reads hide the cleared errors from + the moment the request is recorded. """ try: data = request.get_json(silent=True) or {} + clear_all = _coerce_to_bool(data.get('all')) raw_max_age = data.get('max_age_hours', 24) # Validate and coerce max_age_hours + max_age_hours = None + if not clear_all: + try: + max_age_hours = int(raw_max_age) + if max_age_hours < 1: + return error_response( + error_code=ErrorCode.INVALID_INPUT, + message="max_age_hours must be at least 1", + context={'provided_value': raw_max_age}, + status_code=400 + ) + if max_age_hours > 8760: # 1 year max + return error_response( + error_code=ErrorCode.INVALID_INPUT, + message="max_age_hours cannot exceed 8760 (1 year)", + context={'provided_value': raw_max_age}, + status_code=400 + ) + except (ValueError, TypeError, OverflowError): + return error_response( + error_code=ErrorCode.INVALID_INPUT, + message="max_age_hours must be a valid integer", + context={'provided_value': str(raw_max_age)}, + status_code=400 + ) + + now = _pkg.time.time() + cutoff = now if clear_all else now - max_age_hours * 3600 try: - max_age_hours = int(raw_max_age) - if max_age_hours < 1: - return error_response( - error_code=ErrorCode.INVALID_INPUT, - message="max_age_hours must be at least 1", - context={'provided_value': raw_max_age}, - status_code=400 - ) - if max_age_hours > 8760: # 1 year max - return error_response( - error_code=ErrorCode.INVALID_INPUT, - message="max_age_hours cannot exceed 8760 (1 year)", - context={'provided_value': raw_max_age}, - status_code=400 - ) - except (ValueError, TypeError, OverflowError): + result = _errors.request_error_clear(_errors_cache(), cutoff) + except OSError as e: + logger.error("Could not record an error clear request: %s", e) return error_response( - error_code=ErrorCode.INVALID_INPUT, - message="max_age_hours must be a valid integer", - context={'provided_value': str(raw_max_age)}, - status_code=400 + error_code=ErrorCode.SYSTEM_ERROR, + message="Could not record the clear request in the shared cache", + status_code=500 ) - aggregator = get_error_aggregator() - cleared_count = aggregator.clear_old_records(max_age_hours=max_age_hours) - + scope = "all errors" if clear_all else f"errors older than {max_age_hours} hours" return success_response( - data={'cleared_count': cleared_count}, - message=f"Cleared {cleared_count} error records older than {max_age_hours} hours" + data=result, + message=(f"Clear of {scope} requested; the display service applies it " + f"within about {int(_errors.SNAPSHOT_TICK_INTERVAL)} seconds") ) except Exception as e: logger.error(f"Error clearing old errors: {e}", exc_info=True) diff --git a/web_interface/templates/v3/partials/logs.html b/web_interface/templates/v3/partials/logs.html index 3a23b8b9..81b3c4fb 100644 --- a/web_interface/templates/v3/partials/logs.html +++ b/web_interface/templates/v3/partials/logs.html @@ -104,6 +104,25 @@
Connected to log stream + + +
+
+
+

Plugin errors

+

Recorded by the display service

+
+
+ + +
+
+
Loading...
+