From 4a1fd7464ad7875e16b8bff5c50814a9f831a87b Mon Sep 17 00:00:00 2001 From: Chuck <33324927+ChuckBuilds@users.noreply.github.com> Date: Wed, 23 Sep 2026 14:32:02 -0400 Subject: [PATCH] fix(errors): serve /api/v3/errors/* from the display service; add a Plugin errors panel (#614) * fix(errors): serve /api/v3/errors/* from the display service's aggregator The error aggregator is a per-process singleton and only the display service runs plugins, so only its aggregator records anything. The routes read the web process's own, empty one and always reported no errors. The display service now publishes a bounded snapshot of its aggregator to the shared cache (plugin_error_snapshot) from a daemon thread: at most once every 10 s and only when something changed, never raising into the caller. The routes read it and keep their response shapes, adding snapshot_available, generated_at and clear_pending; exception text has credentials redacted. POST /errors/clear writes a clear request (plugin_error_clear_request) that the display applies on its next 5 s tick via the new clear_before(), which keeps errors recorded after the cutoff and rebuilds the counts. Until the snapshot acknowledges the request, reads hide everything before the cutoff, so a snapshot written just before the click cannot bring errors back. Adds "all": true; cleared_count is null when only the display can know it. Co-Authored-By: Claude Opus 5.5 * feat(web): show plugin errors in the Logs tab A compact panel under the log viewer: per-plugin error counts, repeating errors (type, count, affected plugins, a sample message, last seen) and a Clear button, with empty states for "no errors" and "display service hasn't reported yet". Polls every 15 s while the tab is active; all text goes through escapeHtml. Co-Authored-By: Claude Opus 5.5 * docs: describe where plugin error reports come from and how clear works Co-Authored-By: Claude Opus 5.5 * fix(errors): redact the published snapshot before clipping it Keeping only a traceback's tail (or clipping a message) could cut an `api_key=` marker off while keeping the secret after it, and the web side's redaction would then have nothing to match. The display now redacts every free-text field of the snapshot first. The patterns move to a Flask-free src/redaction.py so the display service can use them; redact_text in the web error handler uses the same function, unchanged in behaviour. Co-Authored-By: Claude Opus 5.5 --------- Co-authored-by: Claude Opus 5.5 --- CHANGELOG.md | 17 + docs/PLUGIN_ERROR_HANDLING.md | 18 +- docs/REST_API_REFERENCE.md | 84 +++- src/display_controller.py | 3 + src/error_aggregator.py | 470 +++++++++++++++++- src/redaction.py | 49 ++ src/web_interface/error_handler.py | 40 +- test/test_error_snapshot_cross_process.py | 443 +++++++++++++++++ web_interface/blueprints/api_v3/misc.py | 133 +++-- web_interface/templates/v3/partials/logs.html | 149 ++++++ 10 files changed, 1327 insertions(+), 79 deletions(-) create mode 100644 src/redaction.py create mode 100644 test/test_error_snapshot_cross_process.py 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...
+