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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

* docs: describe where plugin error reports come from and how clear works

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Chuck
2026-09-23 14:32:02 -04:00
committed by GitHub
co-authored by Claude Opus 5.5
parent cd5a4e2251
commit 4a1fd7464a
10 changed files with 1327 additions and 79 deletions
+17
View File
@@ -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/<id>` 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):
+16 -2
View File
@@ -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.
+79 -5
View File
@@ -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/<plugin_id>`
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
+3
View File
@@ -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
+469 -1
View File
@@ -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/<id> payload: get_plugin_health()'s shape plus
``generated_at``, ``snapshot_available`` and ``clear_pending``."""
health: Dict[str, Any] = {
"plugin_id": plugin_id,
"status": "healthy",
"total_errors": 0,
"error_types": {},
"recent_error_count": 0,
"last_error": None,
}
cutoff = _pending_cutoff(snapshot, clear_request)
table = snapshot.get("plugin_health") if snapshot else None
entry = table.get(plugin_id) if isinstance(table, dict) else None
if isinstance(entry, dict):
# last_error is the plugin's newest error: if even that predates a
# pending clear, so does everything else the display reported for it.
if cutoff is None or _is_after(entry.get("last_error"), "timestamp", cutoff):
for name in ("status", "total_errors", "error_types", "recent_error_count", "last_error"):
if name in entry:
health[name] = entry[name]
health["generated_at"] = snapshot.get("generated_at") if snapshot else None
health["snapshot_available"] = snapshot is not None
health["clear_pending"] = cutoff is not None
return health
def _count_cleared(summary: Dict[str, Any], cutoff: float) -> Optional[int]:
"""How many of the reported errors a clear at ``cutoff`` hides, if known.
Exact when every reported error predates the cutoff (always the case for
a clear of everything) or when the report lists every error; otherwise
only the display service knows, and None says so.
"""
total = summary["total_errors"]
times = [_epoch(r.get("timestamp")) if isinstance(r, dict) else None
for r in summary["recent_errors"]]
if total == 0 or (times and times[-1] is not None and times[-1] <= cutoff):
return total
if len(times) >= total and None not in times:
return sum(1 for t in times if t <= cutoff)
return None
def request_error_clear(cache_manager: Any, cutoff: float) -> Dict[str, Any]:
"""Ask the display service to forget errors recorded at or before ``cutoff``.
Returns ``request_id``, ``cutoff`` (ISO, local time), ``cleared_count``
(see _count_cleared) and ``clear_requested``. Raises OSError when the
request did not reach the shared cache, since a cache without a usable
directory accepts set() and keeps nothing.
A request the display has not applied yet is only ever widened: a later,
narrower one ("older than 24 hours" after "everything") overwriting it
would otherwise bring back the errors the first one hid.
"""
snapshot, clear_request = read_error_report(cache_manager)
pending = _pending_cutoff(snapshot, clear_request)
if pending is not None:
cutoff = max(cutoff, pending)
before = error_summary_from_report(snapshot, clear_request)
request = {
"request_id": uuid.uuid4().hex,
"cutoff": cutoff,
"requested_at": time.time(),
}
cache_manager.set(ERROR_CLEAR_REQUEST_KEY, request)
stored = cache_manager.get(ERROR_CLEAR_REQUEST_KEY, max_age=None, memory_ttl=0)
if not isinstance(stored, dict) or stored.get("request_id") != request["request_id"]:
raise OSError("the clear request was not stored in the shared cache")
return {
"cleared_count": _count_cleared(before, cutoff),
"clear_requested": True,
"request_id": request["request_id"],
"cutoff": datetime.fromtimestamp(cutoff).isoformat(),
}
+49
View File
@@ -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: <scheme> <credential>`. 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 ``<redacted>``; 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<redacted>\3', text)
text = _REDACT_AUTH_HEADER.sub(r'\1<redacted>', text)
return _REDACT_CREDENTIAL.sub(r'\1<redacted>', text)
+2 -38
View File
@@ -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: <scheme> <credential>`. 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<redacted>\3', text)
text = _REDACT_AUTH_HEADER.sub(r'\1<redacted>', text)
text = _REDACT_CREDENTIAL.sub(r'\1<redacted>', 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:
+443
View File
@@ -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=<redacted>" 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))
+100 -33
View File
@@ -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)
@@ -104,6 +104,25 @@
<div class="w-2 h-2 bg-green-500 rounded-full"></div>
<span>Connected to log stream</span>
</div>
<!-- Plugin errors, as recorded by the display service -->
<div id="plugin-errors-panel" class="mt-6 border-t border-gray-200 pt-4">
<div class="flex flex-wrap items-center justify-between gap-2 mb-3">
<div>
<h3 class="text-base font-semibold text-gray-900">Plugin errors</h3>
<p id="plugin-errors-meta" class="text-xs text-gray-600">Recorded by the display service</p>
</div>
<div class="flex items-center gap-2">
<button id="plugin-errors-refresh-btn" type="button" class="btn bg-gray-600 hover:bg-gray-700 text-white px-3 py-1 rounded text-sm">
<i class="fas fa-sync-alt mr-1"></i>Refresh
</button>
<button id="plugin-errors-clear-btn" type="button" class="btn bg-red-600 hover:bg-red-700 text-white px-3 py-1 rounded text-sm">
<i class="fas fa-eraser mr-1"></i>Clear
</button>
</div>
</div>
<div id="plugin-errors-body" class="text-sm text-gray-600" aria-live="polite">Loading...</div>
</div>
</div>
<script>
@@ -218,10 +237,18 @@ window._filteredLogs = [];
// Current-plugin poll and the realtime log stream run only while the Logs
// tab is active and the page is visible (both restart with an immediate
// refresh). Re-registering on partial reload replaces the old one.
const pluginErrorsRefreshBtn = document.getElementById('plugin-errors-refresh-btn');
if (pluginErrorsRefreshBtn) pluginErrorsRefreshBtn.onclick = refreshPluginErrors;
const pluginErrorsClearBtn = document.getElementById('plugin-errors-clear-btn');
if (pluginErrorsClearBtn) pluginErrorsClearBtn.onclick = clearPluginErrors;
function startLogsActivity() {
refreshCurrentPluginStatus();
if (window._currentPluginPollTimer) clearInterval(window._currentPluginPollTimer);
window._currentPluginPollTimer = setInterval(refreshCurrentPluginStatus, 5000);
refreshPluginErrors();
if (window._pluginErrorsPollTimer) clearInterval(window._pluginErrorsPollTimer);
window._pluginErrorsPollTimer = setInterval(refreshPluginErrors, 15000);
if (window._isRealtime && !window._logsEventSource) setupRealtimeLogs();
}
function stopLogsActivity() {
@@ -229,6 +256,10 @@ window._filteredLogs = [];
clearInterval(window._currentPluginPollTimer);
window._currentPluginPollTimer = null;
}
if (window._pluginErrorsPollTimer) {
clearInterval(window._pluginErrorsPollTimer);
window._pluginErrorsPollTimer = null;
}
stopRealtimeLogs();
}
if (window.LEDVisibility) {
@@ -787,6 +818,120 @@ function refreshCurrentPluginStatus() {
});
}
// Plugin errors panel. The display service records them and publishes a
// snapshot every few seconds; /api/v3/errors/summary serves that snapshot.
// Plugin ids and messages come from plugins: everything goes through escapeHtml.
function formatErrorTime(iso) {
return iso ? String(iso).replace('T', ' ').slice(0, 19) : '';
}
function renderPluginErrors(summary) {
const body = document.getElementById('plugin-errors-body');
const meta = document.getElementById('plugin-errors-meta');
if (!body) return;
if (meta) {
let text = 'Recorded by the display service';
if (summary.generated_at) text += ' · updated ' + formatErrorTime(summary.generated_at);
if (summary.clear_pending) text += ' · clearing…';
meta.textContent = text;
}
if (!summary.snapshot_available) {
body.innerHTML = '<p class="text-gray-600"><i class="fas fa-hourglass-half mr-1"></i>' +
'Display service hasn’t reported yet. Errors appear here once it is running.</p>';
return;
}
const plugins = Object.entries(summary.plugin_error_counts || {}).map(([id, types]) => {
const total = Object.values(types || {}).reduce((sum, n) => sum + (Number(n) || 0), 0);
return { id, types: types || {}, total };
}).sort((a, b) => b.total - a.total);
const patterns = Object.values(summary.active_patterns || {});
if (!Number(summary.total_errors) && plugins.length === 0 && patterns.length === 0) {
body.innerHTML = '<p class="text-gray-600"><i class="fas fa-check-circle text-green-600 mr-1"></i>' +
'No plugin errors recorded</p>';
return;
}
let html = '';
if (plugins.length) {
html += '<div class="overflow-x-auto"><table class="min-w-full text-sm">' +
'<thead><tr class="text-left text-xs uppercase text-gray-500 border-b border-gray-200">' +
'<th class="py-1 pr-4 font-medium">Plugin</th><th class="py-1 pr-4 font-medium">Errors</th>' +
'<th class="py-1 font-medium">Types</th></tr></thead><tbody>';
plugins.forEach(p => {
const types = Object.entries(p.types)
.map(([type, n]) => `${escapeHtml(type)} ×${escapeHtml(String(Number(n) || 0))}`)
.join(', ');
html += '<tr class="border-b border-gray-100">' +
`<td class="py-1 pr-4"><code class="bg-gray-100 px-1 rounded">${escapeHtml(p.id)}</code></td>` +
`<td class="py-1 pr-4 font-semibold text-gray-900">${escapeHtml(String(p.total))}</td>` +
`<td class="py-1 text-gray-700 break-words">${types}</td></tr>`;
});
html += '</tbody></table></div>';
}
if (patterns.length) {
const badge = {
critical: 'bg-red-100 text-red-800',
error: 'bg-amber-100 text-amber-800',
warning: 'bg-yellow-100 text-yellow-800'
};
html += '<h4 class="mt-4 mb-2 text-sm font-semibold text-gray-900">Repeating errors</h4><ul class="list-none pl-0 space-y-2">';
patterns.forEach(p => {
const severity = String(p.severity || 'warning');
const affected = (p.affected_plugins || []).map(id => escapeHtml(id)).join(', ') || 'unknown';
const sample = (p.sample_messages || [])[0];
html += '<li class="border border-gray-200 rounded-lg px-3 py-2">' +
'<div class="flex flex-wrap items-center gap-2">' +
`<span class="px-2 py-0.5 rounded text-xs font-semibold ${badge[severity] || badge.warning}">${escapeHtml(severity)}</span>` +
`<code class="font-semibold text-gray-900">${escapeHtml(p.error_type || '')}</code>` +
`<span class="text-gray-700">×${escapeHtml(String(Number(p.count) || 0))}</span>` +
`<span class="text-xs text-gray-500 ml-auto">last seen ${escapeHtml(formatErrorTime(p.last_seen))}</span></div>` +
`<div class="mt-1 text-xs text-gray-600">Plugins: ${affected}</div>` +
(sample ? `<div class="mt-1 text-xs font-mono text-gray-700 break-words">${escapeHtml(sample)}</div>` : '') +
'</li>';
});
html += '</ul>';
}
body.innerHTML = html;
}
function refreshPluginErrors() {
fetch('/api/v3/errors/summary')
.then(response => response.json())
.then(data => {
if (data.status !== 'success' || !data.data) throw new Error(data.message || 'Request failed');
renderPluginErrors(data.data);
})
.catch(() => {
const body = document.getElementById('plugin-errors-body');
if (body) body.textContent = 'Could not load plugin errors.';
});
}
function clearPluginErrors() {
if (!confirm('Clear all recorded plugin errors?')) return;
fetch('/api/v3/errors/clear', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ all: true })
})
.then(response => response.json())
.then(data => {
if (typeof showNotification !== 'undefined') {
showNotification(data.status === 'success' ? 'Plugin errors cleared'
: (data.message || 'Could not clear plugin errors'),
data.status === 'success' ? 'success' : 'error');
}
refreshPluginErrors();
})
.catch(() => {
if (typeof showNotification !== 'undefined') showNotification('Could not clear plugin errors', 'error');
});
}
// Cleanup on page unload
window.addEventListener('beforeunload', function() {
if (window._logsEventSource) {
@@ -796,5 +941,9 @@ window.addEventListener('beforeunload', function() {
clearInterval(window._currentPluginPollTimer);
window._currentPluginPollTimer = null;
}
if (window._pluginErrorsPollTimer) {
clearInterval(window._pluginErrorsPollTimer);
window._pluginErrorsPollTimer = null;
}
});
</script>