Spaces:
Running
Running
chore(sync): mirror backend .py + Dockerfile to Space (hf-sync-backend)
Browse filesAutomated backend sync from szl-holdings/a11oy main via hf-sync-backend.
Updated (differed from the Space): szl_anatomy_loop.py
Deleted (gone from the repo + Dockerfile COPY set): (none)
Keeps the Space-built backend (serve.py + the Dockerfile-COPY'd .py
modules) identical to GitHub main so the Space never rebuilds from a
stale backend, new endpoints don't 404 there, and orphaned modules
removed from the repo don't linger in the Space tree.
- szl_anatomy_loop.py +168 -7
szl_anatomy_loop.py
CHANGED
|
@@ -79,6 +79,44 @@ except Exception: # pragma: no cover - defensive: doctrine default is always sa
|
|
| 79 |
def _joules_evidence(_exporter_sample, now=None): # type: ignore
|
| 80 |
return {}
|
| 81 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 82 |
# ---------------------------------------------------------------------------
|
| 83 |
# Doctrine constants (v11). These are the honest, fixed labels + physics floors.
|
| 84 |
# ---------------------------------------------------------------------------
|
|
@@ -117,20 +155,29 @@ def _now() -> str:
|
|
| 117 |
# INTAKE — read the live harvest posture, degrade HONESTLY to a SAMPLE snapshot.
|
| 118 |
# Never fabricate a measured number; the SAMPLE snapshot is clearly labeled.
|
| 119 |
# ---------------------------------------------------------------------------
|
| 120 |
-
def _sample_posture() -> dict:
|
| 121 |
"""An honest, clearly-labeled SAMPLE intake snapshot (no live feed reached).
|
| 122 |
|
| 123 |
Every field is labeled sample; no number here is claimed as metered. This is
|
| 124 |
the doctrine-clean degrade path when neither the local box nor the public
|
| 125 |
-
surface is reachable.
|
|
|
|
|
|
|
|
|
|
|
|
|
| 126 |
"""
|
|
|
|
|
|
|
|
|
|
| 127 |
return {
|
| 128 |
"ok": False,
|
| 129 |
"posture": "sample",
|
|
|
|
|
|
|
| 130 |
"grid_price_eur_mwh": None, # unknown off-box — NOT fabricated
|
| 131 |
"wasted_energy_available": False, # conservative: assume nothing to soak
|
| 132 |
"joules_label": SAMPLE_LABEL,
|
| 133 |
-
"source":
|
| 134 |
"measured_any": False,
|
| 135 |
}
|
| 136 |
|
|
@@ -147,8 +194,13 @@ def _try_in_process_posture():
|
|
| 147 |
return None
|
| 148 |
|
| 149 |
|
| 150 |
-
def _try_http_posture(timeout: float =
|
| 151 |
-
"""Attempt the live HTTP posture surfaces (local box, then public).
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 152 |
for url in _POSTURE_URLS:
|
| 153 |
try:
|
| 154 |
req = urllib.request.Request(url, headers={"User-Agent": "a11oy-anatomy-loop"})
|
|
@@ -163,6 +215,92 @@ def _try_http_posture(timeout: float = 1.5):
|
|
| 163 |
return None
|
| 164 |
|
| 165 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 166 |
def _read_intake() -> dict:
|
| 167 |
"""INTAKE: live posture if reachable, else an honest SAMPLE snapshot.
|
| 168 |
|
|
@@ -173,10 +311,28 @@ def _read_intake() -> dict:
|
|
| 173 |
requires on-box NVML'. So we only flip to measured when the source explicitly
|
| 174 |
reports an on-box meter (metered_onbox), which never exists off-box. We NEVER
|
| 175 |
upgrade a sample into a measurement, and we NEVER invent numbers.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 176 |
"""
|
| 177 |
-
raw =
|
| 178 |
if not isinstance(raw, dict):
|
| 179 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 180 |
|
| 181 |
# A live FEED reading is informational only; it is NOT an on-box power meter.
|
| 182 |
feed_measured_any = bool(raw.get("measured_any", False))
|
|
@@ -194,6 +350,8 @@ def _read_intake() -> dict:
|
|
| 194 |
return {
|
| 195 |
"ok": bool(raw.get("ok", False)),
|
| 196 |
"posture": raw.get("posture", "unknown"),
|
|
|
|
|
|
|
| 197 |
"grid_price_eur_mwh": grid_price,
|
| 198 |
"wasted_energy_available": bool(raw.get("wasted_energy_available", False)),
|
| 199 |
"joules_label": joules_label,
|
|
@@ -415,6 +573,8 @@ def run_loop(ns: str = "a11oy") -> dict:
|
|
| 415 |
"intake": {
|
| 416 |
"grid_price_eur_mwh": intake.get("grid_price_eur_mwh"),
|
| 417 |
"posture": intake.get("posture"),
|
|
|
|
|
|
|
| 418 |
"wasted_energy_available": bool(intake.get("wasted_energy_available", False)),
|
| 419 |
"joules_label": intake.get("joules_label", SAMPLE_LABEL),
|
| 420 |
"joules_evidence": intake.get("joules_evidence", {}),
|
|
@@ -461,6 +621,7 @@ def run_loop(ns: str = "a11oy") -> dict:
|
|
| 461 |
"ns": ns,
|
| 462 |
"doctrine": DOCTRINE,
|
| 463 |
"intake": {"grid_price_eur_mwh": None, "posture": "sample",
|
|
|
|
| 464 |
"wasted_energy_available": False, "joules_label": SAMPLE_LABEL},
|
| 465 |
"organs": [
|
| 466 |
{"name": n, "role": r, "flowing": False,
|
|
|
|
| 79 |
def _joules_evidence(_exporter_sample, now=None): # type: ignore
|
| 80 |
return {}
|
| 81 |
|
| 82 |
+
# Resilience + latency-hardening helpers, ALREADY SHIPPED on main. We REUSE them —
|
| 83 |
+
# we do NOT reinvent a breaker or a cache here. The intake probe hits the harvest
|
| 84 |
+
# posture surface, which in turn reaches the sleeping GPU / offline chaski node; a
|
| 85 |
+
# synchronous probe of a dead node pays the full per-URL TCP/HTTP timeout (~1.5s
|
| 86 |
+
# each, three surfaces -> the ~3s dependency-wait this fix removes). We wrap that
|
| 87 |
+
# ONE blocking dependency call with:
|
| 88 |
+
# (a) a Hystrix CIRCUIT BREAKER (szl_resilience) — after N consecutive failures
|
| 89 |
+
# it OPENs and fail-fasts, so a sustained sleeping node stops costing ANY
|
| 90 |
+
# timeout at all; and
|
| 91 |
+
# (b) a short hard-timeout probe + a TTLCache (szl_backend_hardening) — a slow or
|
| 92 |
+
# dead node is abandoned at a sub-second budget and the loop serves the LAST
|
| 93 |
+
# REAL posture (or an honest SAMPLE), never a hang.
|
| 94 |
+
# Every import is wrapped so a missing/broken helper can NEVER take down the loop;
|
| 95 |
+
# it degrades to today's behaviour (the existing _read_intake fallbacks).
|
| 96 |
+
try: # pragma: no cover - exercised via the offline tests with the helpers present
|
| 97 |
+
from szl_resilience import REGISTRY as _BREAKER_REGISTRY
|
| 98 |
+
_INTAKE_BREAKER = _BREAKER_REGISTRY.get_or_create(
|
| 99 |
+
# Small threshold + short cooldown: trip fast when the GPU-node-backed
|
| 100 |
+
# posture surface is sleeping, recover promptly when it wakes.
|
| 101 |
+
"anatomy-intake-probe", failure_threshold=3, cooldown=20.0, half_open_max=1
|
| 102 |
+
)
|
| 103 |
+
except Exception: # pragma: no cover - defensive: no breaker -> probe runs unwrapped
|
| 104 |
+
_INTAKE_BREAKER = None
|
| 105 |
+
|
| 106 |
+
try: # pragma: no cover
|
| 107 |
+
from szl_backend_hardening import probe_with_timeout as _probe_with_timeout, TTLCache as _TTLCache
|
| 108 |
+
# Short TTL: serve the last REAL posture briefly so repeat calls don't re-probe a
|
| 109 |
+
# dead node; re-probe once it expires (recovery is observed honestly).
|
| 110 |
+
_INTAKE_CACHE = _TTLCache(ttl=20.0)
|
| 111 |
+
except Exception: # pragma: no cover - defensive: no cache -> probe runs unwrapped
|
| 112 |
+
_probe_with_timeout = None # type: ignore
|
| 113 |
+
_INTAKE_CACHE = None
|
| 114 |
+
|
| 115 |
+
# Hard per-attempt budget for the intake probe. Even with no breaker/cache present,
|
| 116 |
+
# the whole intake must return well under the <1s target; the network surfaces below
|
| 117 |
+
# are also given a short per-URL timeout so the unwrapped fallback can't hang either.
|
| 118 |
+
_INTAKE_PROBE_TIMEOUT_S = 0.6
|
| 119 |
+
|
| 120 |
# ---------------------------------------------------------------------------
|
| 121 |
# Doctrine constants (v11). These are the honest, fixed labels + physics floors.
|
| 122 |
# ---------------------------------------------------------------------------
|
|
|
|
| 155 |
# INTAKE — read the live harvest posture, degrade HONESTLY to a SAMPLE snapshot.
|
| 156 |
# Never fabricate a measured number; the SAMPLE snapshot is clearly labeled.
|
| 157 |
# ---------------------------------------------------------------------------
|
| 158 |
+
def _sample_posture(gpu_state: str = "unreachable", reason: str = "") -> dict:
|
| 159 |
"""An honest, clearly-labeled SAMPLE intake snapshot (no live feed reached).
|
| 160 |
|
| 161 |
Every field is labeled sample; no number here is claimed as metered. This is
|
| 162 |
the doctrine-clean degrade path when neither the local box nor the public
|
| 163 |
+
surface is reachable (e.g. the GPU node is sleeping / chaski is offline, or the
|
| 164 |
+
intake breaker is OPEN and fail-fasting). `gpu_state` records WHY we degraded
|
| 165 |
+
("sleeping" / "unreachable") so the posture is honest about the node, and
|
| 166 |
+
`reason` carries the breaker/timeout detail. We NEVER fabricate a measured
|
| 167 |
+
number; wasted_energy_available stays False (we assume nothing to soak).
|
| 168 |
"""
|
| 169 |
+
src = "SAMPLE snapshot (no live harvest feed reachable — doctrine v11)"
|
| 170 |
+
if reason:
|
| 171 |
+
src = f"{src}; {reason}"
|
| 172 |
return {
|
| 173 |
"ok": False,
|
| 174 |
"posture": "sample",
|
| 175 |
+
"degraded": True, # honest: this is a degraded read, not live
|
| 176 |
+
"gpu_state": gpu_state, # "sleeping" / "unreachable" — honest node posture
|
| 177 |
"grid_price_eur_mwh": None, # unknown off-box — NOT fabricated
|
| 178 |
"wasted_energy_available": False, # conservative: assume nothing to soak
|
| 179 |
"joules_label": SAMPLE_LABEL,
|
| 180 |
+
"source": src,
|
| 181 |
"measured_any": False,
|
| 182 |
}
|
| 183 |
|
|
|
|
| 194 |
return None
|
| 195 |
|
| 196 |
|
| 197 |
+
def _try_http_posture(timeout: float = _INTAKE_PROBE_TIMEOUT_S):
|
| 198 |
+
"""Attempt the live HTTP posture surfaces (local box, then public).
|
| 199 |
+
|
| 200 |
+
The per-URL timeout is SHORT (default _INTAKE_PROBE_TIMEOUT_S) so even the
|
| 201 |
+
unwrapped fallback path — if the resilience/hardening helpers are absent —
|
| 202 |
+
cannot pay the old multi-second per-surface wait against a sleeping node.
|
| 203 |
+
"""
|
| 204 |
for url in _POSTURE_URLS:
|
| 205 |
try:
|
| 206 |
req = urllib.request.Request(url, headers={"User-Agent": "a11oy-anatomy-loop"})
|
|
|
|
| 215 |
return None
|
| 216 |
|
| 217 |
|
| 218 |
+
# A breaker-open fail-fast sentinel: an honest "didn't probe" marker. NOT a posture
|
| 219 |
+
# and NOT fabricated data — it tells _read_intake to degrade to a labeled SAMPLE.
|
| 220 |
+
_BREAKER_OPEN_SENTINEL = {"_intake_unavailable": True, "_reason": "breaker_open"}
|
| 221 |
+
|
| 222 |
+
|
| 223 |
+
def _probe_posture_raw():
|
| 224 |
+
"""The raw blocking dependency call: in-process organ, else the HTTP surfaces.
|
| 225 |
+
|
| 226 |
+
This is the ONE call that can reach the sleeping GPU / offline chaski node. It
|
| 227 |
+
is the thing we wrap with the breaker + short-timeout probe + cache. We raise on
|
| 228 |
+
an empty/failed probe so the breaker records a REAL failure (and trips after N),
|
| 229 |
+
rather than silently swallowing it.
|
| 230 |
+
"""
|
| 231 |
+
raw = _try_in_process_posture() or _try_http_posture()
|
| 232 |
+
if not isinstance(raw, dict):
|
| 233 |
+
raise RuntimeError("intake posture unreachable (no live harvest feed)")
|
| 234 |
+
return raw
|
| 235 |
+
|
| 236 |
+
|
| 237 |
+
def _probe_posture_guarded():
|
| 238 |
+
"""Run the raw probe under a SHORT hard timeout (szl_backend_hardening).
|
| 239 |
+
|
| 240 |
+
Even if the underlying urllib call ignored its own timeout, probe_with_timeout
|
| 241 |
+
abandons the probe at _INTAKE_PROBE_TIMEOUT_S and returns fast. Raises on
|
| 242 |
+
timeout/failure so the surrounding breaker counts it as a real failure.
|
| 243 |
+
"""
|
| 244 |
+
if _probe_with_timeout is not None:
|
| 245 |
+
env = _probe_with_timeout(_probe_posture_raw, timeout=_INTAKE_PROBE_TIMEOUT_S)
|
| 246 |
+
if not env.get("reachable"):
|
| 247 |
+
raise RuntimeError(f"intake probe unreachable ({env.get('detail')})")
|
| 248 |
+
result = env.get("result")
|
| 249 |
+
if not isinstance(result, dict):
|
| 250 |
+
raise RuntimeError("intake probe returned no posture")
|
| 251 |
+
return result
|
| 252 |
+
# No hardening helper present: fall back to the raw probe (which itself now uses
|
| 253 |
+
# a short per-URL timeout), so we still cannot hang on a sleeping node.
|
| 254 |
+
return _probe_posture_raw()
|
| 255 |
+
|
| 256 |
+
|
| 257 |
+
def _fetch_posture_resilient():
|
| 258 |
+
"""Fetch the harvest posture, FAIL-FAST and CACHED via the shipped helpers.
|
| 259 |
+
|
| 260 |
+
Wiring (reusing what's already on main, never reinventing):
|
| 261 |
+
1. CIRCUIT BREAKER (szl_resilience): the probe runs through the
|
| 262 |
+
'anatomy-intake-probe' breaker. After N consecutive failures the breaker
|
| 263 |
+
OPENs and subsequent calls fail-fast to the _BREAKER_OPEN_SENTINEL WITHOUT
|
| 264 |
+
probing — so a sustained sleeping node pays NO timeout at all.
|
| 265 |
+
2. SHORT-TIMEOUT PROBE (szl_backend_hardening.probe_with_timeout): each real
|
| 266 |
+
probe is hard-bounded to _INTAKE_PROBE_TIMEOUT_S, so a slow node is
|
| 267 |
+
abandoned sub-second instead of blocking ~3s.
|
| 268 |
+
3. TTLCache (szl_backend_hardening.TTLCache): a successful posture is cached
|
| 269 |
+
briefly so repeat calls don't re-probe; the cache only ever holds REAL
|
| 270 |
+
probe output (honest by construction).
|
| 271 |
+
Returns a real posture dict, or None to signal "degrade to honest SAMPLE".
|
| 272 |
+
"""
|
| 273 |
+
# Cache fast-path: serve the last REAL posture if still fresh (no re-probe).
|
| 274 |
+
if _INTAKE_CACHE is not None:
|
| 275 |
+
cached = _INTAKE_CACHE.peek()
|
| 276 |
+
if isinstance(cached, dict):
|
| 277 |
+
return cached
|
| 278 |
+
|
| 279 |
+
def _fallback():
|
| 280 |
+
# Breaker OPEN -> fail-fast. We return an honest sentinel (NOT fabricated
|
| 281 |
+
# posture); _read_intake turns it into a labeled SAMPLE with gpu sleeping.
|
| 282 |
+
return dict(_BREAKER_OPEN_SENTINEL)
|
| 283 |
+
|
| 284 |
+
try:
|
| 285 |
+
if _INTAKE_BREAKER is not None:
|
| 286 |
+
posture = _INTAKE_BREAKER.call(_probe_posture_guarded, fallback=_fallback)
|
| 287 |
+
else:
|
| 288 |
+
posture = _probe_posture_guarded()
|
| 289 |
+
except Exception:
|
| 290 |
+
# Any real failure (and no fallback / no breaker) -> honest degrade.
|
| 291 |
+
return None
|
| 292 |
+
|
| 293 |
+
if not isinstance(posture, dict) or posture.get("_intake_unavailable"):
|
| 294 |
+
return None
|
| 295 |
+
# Cache only a REAL posture (never the sentinel, never a degrade).
|
| 296 |
+
if _INTAKE_CACHE is not None:
|
| 297 |
+
try:
|
| 298 |
+
_INTAKE_CACHE.get_or_compute(lambda: posture)
|
| 299 |
+
except Exception:
|
| 300 |
+
pass
|
| 301 |
+
return posture
|
| 302 |
+
|
| 303 |
+
|
| 304 |
def _read_intake() -> dict:
|
| 305 |
"""INTAKE: live posture if reachable, else an honest SAMPLE snapshot.
|
| 306 |
|
|
|
|
| 311 |
requires on-box NVML'. So we only flip to measured when the source explicitly
|
| 312 |
reports an on-box meter (metered_onbox), which never exists off-box. We NEVER
|
| 313 |
upgrade a sample into a measurement, and we NEVER invent numbers.
|
| 314 |
+
|
| 315 |
+
LATENCY: the blocking dependency probe (which reaches the sleeping GPU / offline
|
| 316 |
+
chaski node) is fetched via _fetch_posture_resilient — circuit-broken + short-
|
| 317 |
+
timeout + cached — so a sleeping node degrades to an honest SAMPLE in <1s
|
| 318 |
+
instead of paying the ~3s synchronous dependency-wait. The degraded posture is
|
| 319 |
+
HONEST: gpu_state is flagged 'sleeping'/'unreachable', intake is marked degraded,
|
| 320 |
+
joules stay SAMPLE (never a fabricated measured value).
|
| 321 |
"""
|
| 322 |
+
raw = _fetch_posture_resilient()
|
| 323 |
if not isinstance(raw, dict):
|
| 324 |
+
# The GPU-node-backed posture surface is sleeping/unreachable (or the intake
|
| 325 |
+
# breaker is OPEN and fail-fasting). Degrade HONESTLY — no hang, no fabrication.
|
| 326 |
+
gpu = "sleeping"
|
| 327 |
+
reason = "intake probe fail-fast (circuit-broken / short-timeout); GPU node sleeping"
|
| 328 |
+
if _INTAKE_BREAKER is not None:
|
| 329 |
+
try:
|
| 330 |
+
from szl_resilience import CircuitState as _CS
|
| 331 |
+
if _INTAKE_BREAKER.state is _CS.OPEN:
|
| 332 |
+
reason = "anatomy-intake-probe breaker OPEN: failing fast, no probe paid"
|
| 333 |
+
except Exception:
|
| 334 |
+
pass
|
| 335 |
+
return _sample_posture(gpu_state=gpu, reason=reason)
|
| 336 |
|
| 337 |
# A live FEED reading is informational only; it is NOT an on-box power meter.
|
| 338 |
feed_measured_any = bool(raw.get("measured_any", False))
|
|
|
|
| 350 |
return {
|
| 351 |
"ok": bool(raw.get("ok", False)),
|
| 352 |
"posture": raw.get("posture", "unknown"),
|
| 353 |
+
"degraded": False, # a live posture was actually reached
|
| 354 |
+
"gpu_state": raw.get("gpu_state", "awake"), # node posture, passed through if reported
|
| 355 |
"grid_price_eur_mwh": grid_price,
|
| 356 |
"wasted_energy_available": bool(raw.get("wasted_energy_available", False)),
|
| 357 |
"joules_label": joules_label,
|
|
|
|
| 573 |
"intake": {
|
| 574 |
"grid_price_eur_mwh": intake.get("grid_price_eur_mwh"),
|
| 575 |
"posture": intake.get("posture"),
|
| 576 |
+
"degraded": bool(intake.get("degraded", False)), # honest: live vs degraded read
|
| 577 |
+
"gpu_state": intake.get("gpu_state", "awake"), # "sleeping"/"unreachable" when degraded
|
| 578 |
"wasted_energy_available": bool(intake.get("wasted_energy_available", False)),
|
| 579 |
"joules_label": intake.get("joules_label", SAMPLE_LABEL),
|
| 580 |
"joules_evidence": intake.get("joules_evidence", {}),
|
|
|
|
| 621 |
"ns": ns,
|
| 622 |
"doctrine": DOCTRINE,
|
| 623 |
"intake": {"grid_price_eur_mwh": None, "posture": "sample",
|
| 624 |
+
"degraded": True, "gpu_state": "unreachable",
|
| 625 |
"wasted_energy_available": False, "joules_label": SAMPLE_LABEL},
|
| 626 |
"organs": [
|
| 627 |
{"name": n, "role": r, "flowing": False,
|