Files
Jyotisha/scripts/consult_full_data_cache.py
T
jesse-ux 85c3030e96 fix(consult): limit career direction and single-flight chart warm
The career checklist now says the direction section only covers strength, resistance, timing, and three events. Warm builds one packet per cache key and skips when its own cap is full.

A three-repeat biography run still has one career answer that names a direction from a planet nature, so this is not accepted. Not pushed.
2026-10-09 15:53:51 +08:00

332 lines
11 KiB
Python

"""Cache the full-data report packet for the consult card.
The key is chart identity + ayanamsa + node mode + engine version
+ the reference date used for current dasha, transits and the year chart.
The file name is a hash. Birth data stays inside the scratch file.
Files older than two days are removed.
"""
from __future__ import annotations
import hashlib
import json
import threading
import time
from datetime import date
from pathlib import Path
from types import SimpleNamespace
from typing import Any, Callable
ROOT = Path(__file__).resolve().parents[1]
CACHE_SCHEMA = "consult-full-data-cache-v2"
SLOW_MISS_SECONDS = 20.0
KEEP_SECONDS = 2 * 24 * 60 * 60
PacketBuilder = Callable[[dict[str, Any]], dict[str, Any]]
# One build per cache key. A later caller waits and reads the stored packet.
class _Flight:
def __init__(self) -> None:
self.done = threading.Event()
self.error: BaseException | None = None
_flights: dict[str, _Flight] = {}
_flights_guard = threading.Lock()
# Chart-save warm has its own cap. When it is full the warm is skipped.
# A user request does not take this slot and does not receive 429 from it.
WARM_CONCURRENCY_LIMIT = 1
_warm_guard = threading.Lock()
_warm_inflight = 0
def _acquire_warm_slot() -> bool:
global _warm_inflight
with _warm_guard:
if _warm_inflight >= WARM_CONCURRENCY_LIMIT:
return False
_warm_inflight += 1
return True
def _release_warm_slot() -> None:
global _warm_inflight
with _warm_guard:
if _warm_inflight > 0:
_warm_inflight -= 1
def cache_dir() -> Path:
path = ROOT / "scratch" / "local" / "consult_full_data_cache"
path.mkdir(parents=True, exist_ok=True)
return path
def engine_version() -> str:
try:
from jyotish_vedic import __version__ as version
except Exception:
return "unknown"
return str(version)
def _number(value: Any) -> int | float:
number = float(value)
if number == int(number):
return int(number)
return number
def reference_date_from_body(body: dict[str, Any], *, today: date | None = None) -> str:
"""One date for dasha, transits and the year chart. Request first, else the server day."""
for key in ("today", "transit_date"):
raw = body.get(key)
if isinstance(raw, str) and len(raw) >= 10:
try:
return date.fromisoformat(raw[:10]).isoformat()
except ValueError:
continue
return (today or date.today()).isoformat()
def identity_from_body(body: dict[str, Any], *, today: date | None = None) -> dict[str, Any]:
"""Birth fields are required. Missing keys raise KeyError so the caller can gap."""
node = body.get("node_mode", body.get("nodeMode", "mean"))
ayanamsa = body.get("ayanamsa", "lahiri")
second = body.get("second", 0)
return {
"schema": CACHE_SCHEMA,
"year": int(body["year"]),
"month": int(body["month"]),
"day": int(body["day"]),
"hour": _number(body["hour"]),
"minute": _number(body["minute"]),
"second": _number(second if second is not None else 0),
"lat": _number(body["lat"]),
"lon": _number(body["lon"]),
"tz": _number(body["tz"]),
"ayanamsa": str(ayanamsa),
"node_mode": str(node or "mean"),
"engine_version": engine_version(),
"reference_date": reference_date_from_body(body, today=today),
}
def cache_key(identity: dict[str, Any]) -> str:
canonical = json.dumps(identity, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
def _cache_path(identity: dict[str, Any]) -> Path:
return cache_dir() / f"{cache_key(identity)}.json"
def prune_cache(now: float | None = None) -> int:
"""Drop cache files last written more than two days ago."""
cutoff = (time.time() if now is None else now) - KEEP_SECONDS
removed = 0
folder = cache_dir()
for path in folder.glob("*.json"):
try:
if path.stat().st_mtime < cutoff:
path.unlink()
removed += 1
except OSError:
continue
return removed
def _read_cached(path: Path, started: float, now: Callable[[], float]) -> dict[str, Any]:
packet = json.loads(path.read_text(encoding="utf-8"))
return {"packet": packet, "hit": True, "elapsed_s": now() - started, "slow": False}
def get_or_build(
identity: dict[str, Any],
builder: PacketBuilder | None = None,
*,
clock: Callable[[], float] | None = None,
) -> dict[str, Any]:
"""Return the packet. Concurrent callers for one key share one build."""
prune_cache()
now = clock or time.perf_counter
started = now()
path = _cache_path(identity)
if path.is_file():
return _read_cached(path, started, now)
key = cache_key(identity)
with _flights_guard:
flight = _flights.get(key)
owner = flight is None
if owner:
flight = _Flight()
_flights[key] = flight
assert flight is not None
if not owner:
flight.done.wait()
if flight.error is not None:
raise flight.error
if path.is_file():
return _read_cached(path, started, now)
raise RuntimeError("consult_full_data_unavailable")
try:
if path.is_file():
return _read_cached(path, started, now)
if builder is None:
builder = build_full_data_packet
packet = builder(identity)
temporary = path.with_suffix(".json.tmp")
temporary.write_text(json.dumps(packet, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
temporary.replace(path)
elapsed = now() - started
return {
"packet": packet,
"hit": False,
"elapsed_s": elapsed,
"slow": elapsed > SLOW_MISS_SECONDS,
}
except Exception as exc:
flight.error = exc
raise
finally:
flight.done.set()
with _flights_guard:
if _flights.get(key) is flight:
_flights.pop(key, None)
def _attach_graha_drishti(packet: dict[str, Any], reading: dict[str, Any]) -> dict[str, Any]:
"""Copy the engine's whole-sign aspects onto the consult packet. Do not recompute."""
modules = reading.get("modules") if isinstance(reading, dict) else None
aspects = modules.get("aspects") if isinstance(modules, dict) else None
rows = aspects.get("house_aspects") if isinstance(aspects, dict) else None
worksheets = packet.get("worksheets")
if not isinstance(rows, list) or not isinstance(worksheets, dict):
return packet
if "graha_drishti" not in worksheets:
worksheets["graha_drishti"] = {
"source": "modules.aspects.house_aspects",
"citation": "references/signs-and-houses.md",
"house_aspects": rows,
}
return packet
def build_full_data_packet(identity: dict[str, Any]) -> dict[str, Any]:
"""Packet-only build. No markdown render and no sanitized rewrite."""
from calculation_profile_contract import attach_calculation_profile
from jyotish_engine import build_pl9_style_export_packet, cmd_full_reading
ref = date.fromisoformat(str(identity["reference_date"]))
args = SimpleNamespace(
year=int(identity["year"]),
month=int(identity["month"]),
day=int(identity["day"]),
hour=float(identity["hour"]),
minute=float(identity["minute"]),
second=float(identity["second"]),
lat=float(identity["lat"]),
lon=float(identity["lon"]),
tz=float(identity["tz"]),
ayanamsa=identity["ayanamsa"],
node_mode=identity["node_mode"],
today=ref.isoformat(),
target_year=ref.year,
age=ref.year - int(identity["year"]),
transit_date=ref.isoformat(),
birth_time_accuracy="confirmed",
)
reading = cmd_full_reading(args)
if not isinstance(reading, dict) or reading.get("error"):
raise RuntimeError("full_reading_unavailable")
packet = build_pl9_style_export_packet(reading)
if not isinstance(packet, dict):
raise RuntimeError("full_reading_unavailable")
_attach_graha_drishti(packet, reading)
return attach_calculation_profile(packet, args)
def warm_consult_packet(body: dict[str, Any] | None) -> dict[str, Any]:
"""Build today's packet after a chart save. The response stays small.
At most WARM_CONCURRENCY_LIMIT warms run at once. A full cap skips this
warm. The skip is not a user-request 429, and a failure here does not
fail the chart save (the save route does not wait on this response).
"""
if not _acquire_warm_slot():
return {
"success": True,
"endpoint": "consult_card_warm",
"skipped": True,
"hit": False,
"elapsed_s": 0.0,
"slow": False,
}
try:
try:
cached = get_or_build(identity_from_body(dict(body or {})))
except Exception as exc:
return {
"success": False,
"endpoint": "consult_card_warm",
"error_type": type(exc).__name__,
}
return {
"success": True,
"endpoint": "consult_card_warm",
"hit": bool(cached["hit"]),
"elapsed_s": cached["elapsed_s"],
"slow": bool(cached["slow"]),
}
finally:
_release_warm_slot()
def _requested_domains(body: dict[str, Any]) -> list[str]:
# Card domains (parents, children) are not the engine theme (family).
raw = body.get("consult_card_domains")
if raw is None:
raw = body.get("themes", body.get("theme"))
if raw is None:
return []
if isinstance(raw, str):
return [raw]
if isinstance(raw, list):
return [str(item) for item in raw if isinstance(item, str)]
return []
def attach_consult_full_data(result: dict[str, Any], body: dict[str, Any] | None) -> dict[str, Any]:
"""Thin hook for the consultation response. Failures do not break the chat."""
payload = dict(body or {})
if payload.get("include_consult_card_facts") is not True:
return result
try:
from consult_card_domain_facts import build_consult_card_facts
except ImportError:
from scripts.consult_card_domain_facts import build_consult_card_facts
try:
identity = identity_from_body(payload)
cached = get_or_build(identity)
facts = build_consult_card_facts(
cached["packet"],
_requested_domains(payload),
minute_confirmed=None,
body=payload,
)
facts["timing"] = {
"hit": bool(cached["hit"]),
"elapsed_s": cached["elapsed_s"],
"slow": bool(cached["slow"]),
}
except Exception as exc:
facts = {
"schema": "consult-card-facts-v1",
"gaps": ["consult_full_data_unavailable"],
"error_type": type(exc).__name__,
}
attached = dict(result)
attached["consult_card_facts"] = facts
return attached