mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
feat(fetch): one ESPN scoreboard cache key and a max-age response cache (fetch service stage 2) (#728)
Fetch service stage 2: one ESPN scoreboard cache key shared across the sports base classes (legacy keys still read), and a max-age response cache in the fetch service. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+238
-12
@@ -46,9 +46,25 @@ own conditional headers gets the raw answer. (ESPN sent no validators when
|
||||
this was written -- see the PR that added this module -- so on ESPN the store
|
||||
stays empty and costs nothing.)
|
||||
|
||||
**Response cache (stage 2).** A 200 that says ``Cache-Control: max-age=N``
|
||||
is kept in memory for those N seconds (less its ``Age``), and an identical
|
||||
GET inside that window is answered from it without a request. ESPN sends
|
||||
max-age (1-496 s measured on 2026-10-02, most of it under 10 s) and no
|
||||
validators, so this is the only revalidation-free reuse ESPN allows. It
|
||||
never hands a caller a response older than the caller accepts: a caller
|
||||
says how old with ``cache_max_age`` (``fetch_get(..., cache_max_age=ttl)``;
|
||||
0 skips the cache), and one that does not say gets at most
|
||||
``response_cache.default_max_age`` (30 s). ``no-store``, ``no-cache``,
|
||||
``private``, ``Vary: *`` and ``Set-Cookie`` responses are never kept.
|
||||
Identical means what the validator store keys on: URL, query, effective
|
||||
headers and, for a session with cookies or auth, the session.
|
||||
|
||||
**Counters.** Requests, merged requests, bytes, 304s, errors, HTTP errors,
|
||||
adapter retries, throttled requests and seconds waited, per plugin and per
|
||||
host. :class:`FetchStatsPublisher` publishes them for the web interface
|
||||
adapter retries, throttled requests and seconds waited, plus requests
|
||||
answered without the network: ``memo_hits`` (the response cache) and
|
||||
``cache_hits`` / ``legacy_cache_hits`` (a shared ESPN scoreboard cache entry,
|
||||
counted by ``src/common/espn_dates.py``). Per plugin and per host.
|
||||
:class:`FetchStatsPublisher` publishes them for the web interface
|
||||
(``GET /api/v3/plugins/fetch-stats``).
|
||||
|
||||
CALLER IDENTITY
|
||||
@@ -68,7 +84,8 @@ Counters are kept per plugin without plugins saying who they are:
|
||||
3. Otherwise the request is the core's own (``"core"``).
|
||||
|
||||
Nothing here raises on account of bookkeeping: a failure in counting,
|
||||
keying or the validator store falls back to a plain ``session.get``.
|
||||
keying, the validator store or the response cache falls back to a plain
|
||||
``session.get``.
|
||||
|
||||
Core-internal for now (stage 1). Plugins reach it through ``APIHelper`` and
|
||||
``espn_dates``; a plugin-facing API comes with stage 3.
|
||||
@@ -143,8 +160,22 @@ DEFAULT_CONFIG: Mapping[str, Any] = {
|
||||
"max_bytes": 4 * 1024 * 1024,
|
||||
"max_entry_bytes": 1024 * 1024,
|
||||
},
|
||||
# Stage 2: responses ESPN calls fresh (Cache-Control: max-age), reused
|
||||
# for identical GETs. A college-football Saturday is ~1 MB decoded, so
|
||||
# one entry may be 2 MB; months (5-7 MB) are never kept.
|
||||
"response_cache": {
|
||||
"enabled": True,
|
||||
"default_max_age": 30,
|
||||
"max_entries": 64,
|
||||
"max_bytes": 6 * 1024 * 1024,
|
||||
"max_entry_bytes": 2 * 1024 * 1024,
|
||||
},
|
||||
}
|
||||
|
||||
#: However long a server says a response stays fresh, it is not kept longer
|
||||
#: than this: the cache is for requests that coincide, not for storage.
|
||||
_RESPONSE_CACHE_CEILING = 600.0
|
||||
|
||||
#: Connection pools kept per shared adapter (one per host) and connections
|
||||
#: kept per pool. Larger than requests' 10 because one adapter now serves
|
||||
#: every core Session with its retry policy: three background workers each
|
||||
@@ -171,6 +202,9 @@ _COUNTER_FIELDS = (
|
||||
"overruns", # requests that went after max_wait_seconds anyway
|
||||
"bytes", # decoded response body bytes received
|
||||
"wait_seconds", # time spent waiting for host budgets
|
||||
"memo_hits", # answered from the response cache (max-age); nothing sent
|
||||
"cache_hits", # scoreboard fetches answered from a shared ESPN cache entry
|
||||
"legacy_cache_hits", # cache reads answered from a pre-stage-2 key (any helper)
|
||||
)
|
||||
|
||||
|
||||
@@ -393,6 +427,117 @@ class _ValidatorStore:
|
||||
return {"entries": len(self._entries), "bytes": self._bytes}
|
||||
|
||||
|
||||
# --- response cache (Cache-Control: max-age) --------------------------------------
|
||||
|
||||
@dataclass
|
||||
class _Fresh:
|
||||
response: requests.Response
|
||||
stored_at: float
|
||||
#: Seconds after stored_at the server said the response stays fresh.
|
||||
lifetime: float
|
||||
size: int
|
||||
|
||||
|
||||
class _ResponseCache:
|
||||
"""LRU of finished 200 responses, each kept for its server max-age.
|
||||
|
||||
:meth:`get` answers only while the entry is younger than both its own
|
||||
lifetime and the caller's limit, so nobody is handed a response older
|
||||
than they asked for. Expired entries are dropped as they are met and on
|
||||
every insert, so the cache holds only what is still fresh.
|
||||
"""
|
||||
|
||||
def __init__(self, max_entries: int, max_bytes: int, max_entry_bytes: int,
|
||||
clock: Callable[[], float]) -> None:
|
||||
self.max_entries = max_entries
|
||||
self.max_bytes = max_bytes
|
||||
self.max_entry_bytes = max_entry_bytes
|
||||
self._clock = clock
|
||||
self._entries: "OrderedDict[Any, _Fresh]" = OrderedDict()
|
||||
self._bytes = 0
|
||||
self._lock = threading.Lock()
|
||||
|
||||
def get(self, key: Any, max_age: float) -> Optional[requests.Response]:
|
||||
with self._lock:
|
||||
entry = self._entries.get(key)
|
||||
if entry is None:
|
||||
return None
|
||||
age = self._clock() - entry.stored_at
|
||||
if age < 0 or age >= entry.lifetime:
|
||||
self._drop_locked(key)
|
||||
return None
|
||||
if age > max_age:
|
||||
return None # fresh for someone less strict; kept
|
||||
self._entries.move_to_end(key)
|
||||
clone: requests.Response = _clone_response(entry.response)
|
||||
return clone
|
||||
|
||||
def put(self, key: Any, response: requests.Response, lifetime: float,
|
||||
size: int) -> None:
|
||||
with self._lock:
|
||||
self._drop_locked(key)
|
||||
if size > self.max_entry_bytes or self.max_entries <= 0 or lifetime <= 0:
|
||||
return
|
||||
now = self._clock()
|
||||
for old_key in [k for k, e in self._entries.items()
|
||||
if now - e.stored_at >= e.lifetime]:
|
||||
self._drop_locked(old_key)
|
||||
self._entries[key] = _Fresh(_clone_response(response), now, lifetime, size)
|
||||
self._bytes += size
|
||||
while self._entries and (len(self._entries) > self.max_entries
|
||||
or self._bytes > self.max_bytes):
|
||||
_, old = self._entries.popitem(last=False)
|
||||
self._bytes -= old.size
|
||||
|
||||
def _drop_locked(self, key: Any) -> None:
|
||||
old = self._entries.pop(key, None)
|
||||
if old is not None:
|
||||
self._bytes -= old.size
|
||||
|
||||
def clear(self) -> None:
|
||||
with self._lock:
|
||||
self._entries.clear()
|
||||
self._bytes = 0
|
||||
|
||||
def stats(self) -> Dict[str, int]:
|
||||
with self._lock:
|
||||
return {"entries": len(self._entries), "bytes": self._bytes}
|
||||
|
||||
|
||||
def _cache_directives(value: Optional[str]) -> Dict[str, Optional[str]]:
|
||||
directives: Dict[str, Optional[str]] = {}
|
||||
for part in (value or "").split(","):
|
||||
name, _, arg = part.strip().partition("=")
|
||||
if name:
|
||||
directives[name.strip().lower()] = arg.strip().strip('"') if arg else None
|
||||
return directives
|
||||
|
||||
|
||||
def _fresh_for(response: Any) -> Optional[float]:
|
||||
"""Seconds a finished response stays fresh by its own headers, or None
|
||||
when it must not be reused: not a 200 with its body read, ``no-store``,
|
||||
``no-cache``, ``private``, ``Vary: *``, ``Set-Cookie``, or no max-age."""
|
||||
if _status_of(response) != 200 or _body_of(response) is None:
|
||||
return None
|
||||
directives = _cache_directives(_str_header(response, "Cache-Control"))
|
||||
if {"no-store", "no-cache", "private"} & set(directives):
|
||||
return None
|
||||
if (_str_header(response, "Vary") or "").strip() == "*":
|
||||
return None
|
||||
if _str_header(response, "Set-Cookie"):
|
||||
return None
|
||||
try:
|
||||
max_age = int(directives.get("max-age") or "")
|
||||
except ValueError:
|
||||
return None
|
||||
try:
|
||||
age = int(_str_header(response, "Age") or 0)
|
||||
except ValueError:
|
||||
age = 0
|
||||
fresh = float(min(max_age - max(age, 0), _RESPONSE_CACHE_CEILING))
|
||||
return fresh if fresh > 0 else None
|
||||
|
||||
|
||||
# --- helpers ------------------------------------------------------------------------
|
||||
|
||||
def _host_of(url: Any) -> str:
|
||||
@@ -580,6 +725,9 @@ class FetchService:
|
||||
self.max_wait_seconds = 2.0
|
||||
self._rate_limits: Dict[str, Tuple[float, float]] = {}
|
||||
self._validators = _ValidatorStore(0, 0, 0)
|
||||
self.response_cache = True
|
||||
self.default_max_age = 30.0
|
||||
self._fresh = _ResponseCache(0, 0, 0, clock)
|
||||
self._applied: Optional[str] = None
|
||||
self.configure(config)
|
||||
|
||||
@@ -635,6 +783,14 @@ class FetchService:
|
||||
store = merged.get("validator_store")
|
||||
store = store if isinstance(store, Mapping) else {}
|
||||
default_store = DEFAULT_CONFIG["validator_store"]
|
||||
fresh = merged.get("response_cache")
|
||||
if fresh is not None and not isinstance(fresh, Mapping):
|
||||
logger.warning("fetch_service.response_cache is not an object; using the defaults")
|
||||
fresh = fresh if isinstance(fresh, Mapping) else {}
|
||||
default_fresh = DEFAULT_CONFIG["response_cache"]
|
||||
self.response_cache = fresh.get("enabled") is not False
|
||||
self.default_max_age = _as_float(fresh.get("default_max_age"),
|
||||
float(default_fresh["default_max_age"]))
|
||||
with self._lock:
|
||||
self._rate_limits = limits
|
||||
self._buckets.clear()
|
||||
@@ -643,6 +799,12 @@ class FetchService:
|
||||
_as_int(store.get("max_bytes"), default_store["max_bytes"]),
|
||||
_as_int(store.get("max_entry_bytes"), default_store["max_entry_bytes"]),
|
||||
)
|
||||
self._fresh = _ResponseCache(
|
||||
_as_int(fresh.get("max_entries"), default_fresh["max_entries"]),
|
||||
_as_int(fresh.get("max_bytes"), default_fresh["max_bytes"]),
|
||||
_as_int(fresh.get("max_entry_bytes"), default_fresh["max_entry_bytes"]),
|
||||
self._clock,
|
||||
)
|
||||
self.change_count += 1
|
||||
|
||||
def describe_config(self) -> Dict[str, Any]:
|
||||
@@ -655,6 +817,8 @@ class FetchService:
|
||||
"conditional_get": self.conditional_get,
|
||||
"max_wait_seconds": self.max_wait_seconds,
|
||||
"rate_limits": limits,
|
||||
"response_cache": self.response_cache,
|
||||
"default_max_age": self.default_max_age,
|
||||
}
|
||||
|
||||
def _limit_for(self, host: str) -> Optional[Tuple[float, float]]:
|
||||
@@ -723,7 +887,7 @@ class FetchService:
|
||||
# -- requests --
|
||||
|
||||
def get(self, session: Any, url: str, *, share_in_flight: bool = True,
|
||||
**kwargs: Any) -> Any:
|
||||
cache_max_age: Optional[float] = None, **kwargs: Any) -> Any:
|
||||
"""``session.get(url, **kwargs)`` through the service.
|
||||
|
||||
Same return value, same exceptions, and ``session.get`` is called
|
||||
@@ -734,6 +898,11 @@ class FetchService:
|
||||
than joining an identical one in flight -- for a caller that may
|
||||
retry *because* an earlier request hung (BackgroundDataService
|
||||
cancels and replaces a fetch) and must not be handed that one.
|
||||
|
||||
``cache_max_age`` is the oldest response, in seconds, the caller
|
||||
will take from the response cache (its own TTL); 0 always asks the
|
||||
network. None means ``response_cache.default_max_age``. A response
|
||||
is never reused past the max-age its server gave it either.
|
||||
"""
|
||||
transport = session if session is not None else self.session_for(url)
|
||||
if not self.enabled:
|
||||
@@ -744,8 +913,24 @@ class FetchService:
|
||||
logger.debug("fetch_service could not key a request to %s", url, exc_info=True)
|
||||
request = _Request(plugin=self._caller(), host=_host_of(url))
|
||||
|
||||
reusable = (self.response_cache and request.representation is not None
|
||||
and not _has_conditional_headers(kwargs.get("headers")))
|
||||
if reusable:
|
||||
try:
|
||||
limit = self._accepted_age(cache_max_age)
|
||||
cached = self._fresh.get(request.representation, limit) if limit > 0 else None
|
||||
except Exception:
|
||||
logger.debug("fetch_service response cache lookup failed", exc_info=True)
|
||||
cached = None
|
||||
if cached is not None:
|
||||
self._count(request.plugin, request.host, memo_hits=1)
|
||||
return cached
|
||||
|
||||
if request.flight is None or not self.single_flight or not share_in_flight:
|
||||
return self._send_get(transport, url, kwargs, request)
|
||||
response = self._send_get(transport, url, kwargs, request)
|
||||
if reusable:
|
||||
self._remember(request, response)
|
||||
return response
|
||||
|
||||
with self._lock:
|
||||
flight = self._inflight.get(request.flight)
|
||||
@@ -764,6 +949,8 @@ class FetchService:
|
||||
try:
|
||||
response = self._send_get(transport, url, kwargs, request)
|
||||
flight.response = response
|
||||
if reusable:
|
||||
self._remember(request, response)
|
||||
return response
|
||||
except BaseException as error:
|
||||
flight.error = error
|
||||
@@ -790,6 +977,40 @@ class FetchService:
|
||||
except Exception:
|
||||
logger.debug("fetch_service could not count a merged request", exc_info=True)
|
||||
|
||||
def note_cache_hit(self, url: Any = None, *, legacy: bool = False,
|
||||
avoided_request: bool = True,
|
||||
plugin_id: Optional[str] = None) -> None:
|
||||
"""Count a read answered from a shared cache entry instead of the
|
||||
network (``espn_dates``). ``legacy`` marks a read from a key that
|
||||
predates the canonical one; ``avoided_request=False`` counts only
|
||||
that, for a read that was never going to fetch on a miss."""
|
||||
try:
|
||||
self._count(plugin_id or self._caller(), _host_of(url),
|
||||
cache_hits=int(avoided_request), legacy_cache_hits=int(legacy))
|
||||
except Exception:
|
||||
logger.debug("fetch_service could not count a cache hit", exc_info=True)
|
||||
|
||||
def _accepted_age(self, cache_max_age: Any) -> float:
|
||||
"""The oldest cached response this call accepts, in seconds."""
|
||||
if cache_max_age is None:
|
||||
return self.default_max_age
|
||||
if isinstance(cache_max_age, bool) or not isinstance(cache_max_age, (int, float)):
|
||||
return self.default_max_age
|
||||
value = float(cache_max_age)
|
||||
return value if math.isfinite(value) and value > 0 else 0.0
|
||||
|
||||
def _remember(self, request: _Request, response: Any) -> None:
|
||||
"""Keep a finished response for its server max-age. Never raises."""
|
||||
try:
|
||||
lifetime = _fresh_for(response)
|
||||
if lifetime is None:
|
||||
return
|
||||
body = _body_of(response)
|
||||
self._fresh.put(request.representation, response, lifetime,
|
||||
len(body) if body is not None else 0)
|
||||
except Exception:
|
||||
logger.debug("fetch_service could not keep a response", exc_info=True)
|
||||
|
||||
def _caller(self) -> str:
|
||||
return current_plugin_id() or CORE
|
||||
|
||||
@@ -927,10 +1148,11 @@ class FetchService:
|
||||
per_plugin[name] += value
|
||||
per_host[name] += value
|
||||
self._totals[name] += value
|
||||
if changes.get("requests") or changes.get("merged"):
|
||||
asked = int(changes.get("requests", 0) + changes.get("merged", 0)
|
||||
+ changes.get("memo_hits", 0) + changes.get("cache_hits", 0))
|
||||
if asked:
|
||||
hosts = self._plugin_hosts.setdefault(plugin, {})
|
||||
hosts[host] = hosts.get(host, 0) + int(changes.get("requests", 0)
|
||||
+ changes.get("merged", 0))
|
||||
hosts[host] = hosts.get(host, 0) + asked
|
||||
self.change_count += 1
|
||||
|
||||
def reset_counters(self) -> None:
|
||||
@@ -942,12 +1164,14 @@ class FetchService:
|
||||
self.change_count += 1
|
||||
|
||||
def reset(self) -> None:
|
||||
"""Counters, validators, budgets and in-flight table (tests)."""
|
||||
"""Counters, validators, response cache, budgets and in-flight
|
||||
table (tests)."""
|
||||
self.reset_counters()
|
||||
with self._lock:
|
||||
self._buckets.clear()
|
||||
self._inflight.clear()
|
||||
self._validators.clear()
|
||||
self._fresh.clear()
|
||||
|
||||
def snapshot(self) -> Dict[str, Any]:
|
||||
"""Counters since the service started, JSON-ready."""
|
||||
@@ -971,6 +1195,7 @@ class FetchService:
|
||||
"plugins": plugins,
|
||||
"hosts": hosts,
|
||||
"validators": self._validators.stats(),
|
||||
"response_cache": self._fresh.stats(),
|
||||
"config": self.describe_config(),
|
||||
}
|
||||
|
||||
@@ -1005,10 +1230,11 @@ def configure_fetch_service(config: Any) -> FetchService:
|
||||
|
||||
|
||||
def fetch_get(session: Any, url: str, *, share_in_flight: bool = True,
|
||||
**kwargs: Any) -> Any:
|
||||
"""``session.get(url, **kwargs)`` through the process's FetchService."""
|
||||
cache_max_age: Optional[float] = None, **kwargs: Any) -> Any:
|
||||
"""``session.get(url, **kwargs)`` through the process's FetchService.
|
||||
``cache_max_age``: see :meth:`FetchService.get`."""
|
||||
return get_fetch_service().get(session, url, share_in_flight=share_in_flight,
|
||||
**kwargs)
|
||||
cache_max_age=cache_max_age, **kwargs)
|
||||
|
||||
|
||||
def fetch_post(session: Any, url: str, **kwargs: Any) -> Any:
|
||||
|
||||
Reference in New Issue
Block a user