diff --git a/src/background_data_service.py b/src/background_data_service.py index 22989a56..a41aef45 100644 --- a/src/background_data_service.py +++ b/src/background_data_service.py @@ -26,6 +26,7 @@ from enum import Enum from concurrent.futures import ThreadPoolExecutor import pytz from src.cache_manager import CacheManager +from src.common.json_body import response_json from src.common.espn_dates import ( RANGE_RETRY_SECONDS, _note_range_rejected, @@ -389,7 +390,7 @@ class BackgroundDataService: response.raise_for_status() else: response.raise_for_status() - data = response.json() + data = response_json(response) # Validate data structure if not isinstance(data, dict): diff --git a/src/cache/disk_cache.py b/src/cache/disk_cache.py index 87c3d7a9..fb776b56 100644 --- a/src/cache/disk_cache.py +++ b/src/cache/disk_cache.py @@ -7,6 +7,7 @@ Handles persistent disk-based caching with atomic writes and error recovery. import json import math import os +import re import stat import time import tempfile @@ -98,6 +99,40 @@ def _replace_nonfinite(obj: Any) -> Any: # deleted. Both halves are covered by test/test_cache_nonfinite_floats.py. +#: Enough of a record to hold its header: ``{"timestamp":,"ttl":,``. +_HEAD_BYTES = 256 + +#: A record written with its header first (CacheManager.set does). Anything +#: else -- older files with "data" first, records from other writers -- does not +#: match and is parsed in full, as before. +_HEAD_RE = re.compile( + rb'\A\s*\{\s*"timestamp"\s*:\s*(-?[0-9][0-9.eE+-]*)\s*' + rb'(?:,\s*"ttl"\s*:\s*(-?[0-9][0-9.eE+-]*))?\s*[,}]' +) + + +def _stale_from_head(head: bytes, max_age: Optional[int], now: float) -> bool: + """True when a record's header alone shows it has expired. + + Mirrors the expiry rule in DiskCache.get: a per-entry ttl wins over the + caller's max_age, and no limit at all means never stale. False whenever the + header cannot be read, so the full parse decides as it always did. + """ + match = _HEAD_RE.match(head) + if not match: + return False + try: + timestamp = float(match.group(1)) + limit = max_age + if match.group(2) is not None: + ttl = float(match.group(2)) + if ttl >= 0: + limit = ttl + except ValueError: + return False + return limit is not None and (now - timestamp) > limit + + if orjson is not None: # Encoding the cache record dominated the background fetch worker: on a # Pi 4, stdlib json.dumps runs ~12ms per MB and holds the GIL for all of @@ -266,6 +301,14 @@ class DiskCache: try: with self._lock: with open(cache_path, 'rb') as f: + # Decide staleness from the header before paying for the + # parse. A stale read is the common case for the biggest + # records (a season schedule is re-fetched when its cache + # expires), and parsing 53MB to throw it away held the GIL + # for ~1.8s -- a visible freeze on the panel. + if _stale_from_head(f.read(_HEAD_BYTES), max_age, time.time()): + return None + f.seek(0) record = _loads(f.read()) # Determine record timestamp (prefer embedded, else file mtime) diff --git a/src/cache_manager.py b/src/cache_manager.py index 4d80c986..d064a37f 100644 --- a/src/cache_manager.py +++ b/src/cache_manager.py @@ -522,8 +522,9 @@ class CacheManager: def update_cache(self, data_type: str, data: Dict[str, Any]) -> bool: """Update cache with new data.""" cache_data = { + # Header first; see DiskCache's stale check. + 'timestamp': time.time(), 'data': data, - 'timestamp': time.time() } return self.save_cache(data_type, cache_data) @@ -556,12 +557,15 @@ class CacheManager: from the key and is only a fallback for entries that did not say. Omit it to keep that inferred behaviour. """ - cache_data = { - 'data': data, - 'timestamp': time.time() - } + # timestamp and ttl before data, so they are the first bytes on disk: + # DiskCache.get reads them from the head of the file and can call a + # record stale without parsing it. That matters for the big ones -- a + # whole MLB season is 53MB and ~1.8s of orjson.loads with the GIL held, + # paid in full only to learn the record had expired. + cache_data: Dict[str, Any] = {'timestamp': time.time()} if ttl is not None: cache_data['ttl'] = ttl + cache_data['data'] = data self.save_cache(key, cache_data) @deprecated("3.7.0") diff --git a/src/common/espn_dates.py b/src/common/espn_dates.py index 9be705d1..bda7e86a 100644 --- a/src/common/espn_dates.py +++ b/src/common/espn_dates.py @@ -39,6 +39,14 @@ from datetime import date, timedelta from functools import partial from typing import Any, Dict, List, Optional, Tuple +try: + from src.common.json_body import response_json +except ImportError: + # Plugins bundle copies of this module for older cores, which predate + # json_body; the stdlib parse is what those cores always used. + def response_json(response: Any) -> Any: + return response.json() + # Above this, ESPN returns a truncated list instead of an error. See module # docstring: 500 is the largest value measured to return complete data. ESPN_MAX_LIMIT = 500 @@ -194,7 +202,7 @@ def _fetch_one_chunk( timeout=timeout, ) response.raise_for_status() - return response.json() + return response_json(response) except Exception as exc: # noqa: BLE001 - see docstring if logger: logger.warning("ESPN chunk %s failed, skipping it: %s", chunk, exc) @@ -371,4 +379,4 @@ def fetch_espn_scoreboard( if data is not None: return data response.raise_for_status() - return response.json() + return response_json(response) diff --git a/src/common/json_body.py b/src/common/json_body.py new file mode 100644 index 00000000..fb89301d --- /dev/null +++ b/src/common/json_body.py @@ -0,0 +1,29 @@ +"""Parse an HTTP response body as JSON, with orjson when it is installed. + +``requests``' ``response.json()`` uses the stdlib parser. For the payloads the +sports plugins fetch -- a season schedule is tens of MB -- that runs ~1.7x +slower than orjson on a Pi 4 (3.1s against 1.8s for the 53MB MLB season), and +both hold the GIL for the whole parse, which freezes the display for as long. +Nothing else changes: the result is the same Python objects. +""" + +from __future__ import annotations + +from typing import Any + +try: + import orjson +except ImportError: # optional dependency; see docs/SCROLL_PERFORMANCE.md + orjson = None + + +def response_json(response: Any) -> Any: + """``response.json()``, parsed by orjson when available.""" + body = getattr(response, "content", None) + if orjson is None or not isinstance(body, (bytes, bytearray)): + return response.json() + try: + return orjson.loads(body) + except orjson.JSONDecodeError: + # Let requests raise its usual error, with its usual message. + return response.json() diff --git a/test/test_cache_stale_header.py b/test/test_cache_stale_header.py new file mode 100644 index 00000000..026b982d --- /dev/null +++ b/test/test_cache_stale_header.py @@ -0,0 +1,118 @@ +"""A stale cache record is recognised from its header, without parsing it. + +The sports plugins cache whole season schedules -- 53MB for MLB, 18MB for NHL. +When one expired, DiskCache.get parsed all of it (~1.8s of orjson.loads on a +Pi 4, GIL held, the whole display frozen) only to find the timestamp too old +and throw the result away. CacheManager.set now writes timestamp and ttl ahead +of the data, and DiskCache.get reads them from the first bytes of the file. +""" + +import json +import time +from types import SimpleNamespace + +import pytest + +from src.cache import disk_cache as disk_cache_module +from src.cache.disk_cache import DiskCache, _stale_from_head +from src.common import json_body + + +@pytest.fixture +def disk(tmp_path): + return DiskCache(cache_dir=str(tmp_path)) + + +@pytest.fixture +def parses(monkeypatch): + """Count full parses of cache files.""" + calls = [] + real = disk_cache_module._loads + + def counting(raw): + calls.append(len(raw)) + return real(raw) + + monkeypatch.setattr(disk_cache_module, "_loads", counting) + return calls + + +def _header_first(age=0.0, ttl=None, events=100): + record = {"timestamp": time.time() - age} + if ttl is not None: + record["ttl"] = ttl + record["data"] = {"events": [{"id": n, "name": "x" * 50} for n in range(events)]} + return record + + +def test_cache_manager_writes_the_header_first(monkeypatch): + from src.cache_manager import CacheManager + written = {} + manager = CacheManager.__new__(CacheManager) + monkeypatch.setattr(manager, "save_cache", + lambda key, record: written.update({key: record}), + raising=False) + CacheManager.set(manager, "k", {"events": []}, ttl=60) + assert list(written["k"]) == ["timestamp", "ttl", "data"] + CacheManager.set(manager, "k", {"events": []}) + assert list(written["k"]) == ["timestamp", "data"] + + +def test_a_stale_record_is_not_parsed(disk, parses): + disk.set("season", _header_first(age=600)) + assert disk.get("season", max_age=300) is None + assert parses == [] + + +def test_a_fresh_record_is_parsed_and_returned(disk, parses): + disk.set("season", _header_first(age=10)) + record = disk.get("season", max_age=300) + assert record["data"]["events"][0]["id"] == 0 + assert len(parses) == 1 + + +def test_the_entry_ttl_wins_over_max_age(disk, parses): + disk.set("long", _header_first(age=600, ttl=3600)) + assert disk.get("long", max_age=300) is not None # ttl says fresh + disk.set("short", _header_first(age=60, ttl=30)) + parses.clear() + assert disk.get("short", max_age=300) is None # ttl says stale + assert parses == [] + + +def test_no_limit_means_never_stale(disk): + disk.set("forever", _header_first(age=10 ** 7)) + assert disk.get("forever", max_age=None) is not None + + +def test_older_files_with_data_first_still_work(disk, parses): + # Records written before the header moved: parsed in full, as before. + disk.set("legacy_fresh", {"data": {"v": 1}, "timestamp": time.time()}) + disk.set("legacy_stale", {"data": {"v": 1}, "timestamp": time.time() - 600}) + assert disk.get("legacy_fresh", max_age=300)["data"] == {"v": 1} + assert disk.get("legacy_stale", max_age=300) is None + assert len(parses) == 2 + + +@pytest.mark.parametrize("head, stale", [ + (b'{"timestamp":100.0,"data":{}}', True), + (b'{"timestamp": 100.0, "ttl": 1000, "data": {}}', False), # stdlib spacing + (b'{"timestamp":1e2,"ttl":5,"data":1}', True), + (b'{"timestamp":100.0}', True), + (b'{"data":{},"timestamp":100.0}', False), # unknown layout + (b'{"timestamp":"100.0","data":{}}', False), # string: parse it + (b'', False), +]) +def test_reading_the_header(head, stale): + assert _stale_from_head(head, 300, now=1000.0) is stale + + +def test_response_json_prefers_orjson_and_falls_back(): + payload = {"events": [1, 2, 3]} + response = SimpleNamespace(content=json.dumps(payload).encode(), + json=lambda: pytest.fail("used the slow path")) + if json_body.orjson is None: + pytest.skip("orjson not installed") + assert json_body.response_json(response) == payload + # A response object without bytes content (a test double) still works. + assert json_body.response_json(SimpleNamespace(json=lambda: payload)) == payload