"""Consult-only VedAstro snapshot for the rectification gate. The foreground gateway (BUG-301) still runs. This module is the other call, inside the consultation rectification gate, which used to start the same 4-second official snapshot runner and then drop the result. When the request already asked to defer optional external evidence: - reuse a BUG-727 snapshot for the same birth, ayanamsa, node, and UTC date; - otherwise return a negative-cache hit from earlier today; - otherwise do not start the runner. Record ``official_closure_reason=deferred_in_consultation`` until the end of that UTC day. The rectification HTTP route does not set the flag, so it still calls the real gateway. Negative entries live in their own directory and are never served through the 7-day stale window. """ from __future__ import annotations import contextlib import json import os import re import tempfile from datetime import datetime from pathlib import Path from typing import Any try: from scripts.vedastro_snapshot_cache import ( annotate_stale_gateway, lookup_snapshot, official_snapshot_reference_date, snapshot_cache_key, ) except ModuleNotFoundError: # pragma: no cover - script execution path from vedastro_snapshot_cache import ( annotate_stale_gateway, lookup_snapshot, official_snapshot_reference_date, snapshot_cache_key, ) NEGATIVE_CACHE_SCHEMA = "vedastro_snapshot_negative_cache.v1" DEFERRED_IN_CONSULTATION = "deferred_in_consultation" _SHA256_NAME = re.compile(r"^[0-9a-f]{64}\.json$") _IDENTITY_FIELDS = ( "year", "month", "day", "hour", "minute", "second", "lat", "lon", "tz", "ayanamsa", "ayanamsa_name", "ayanamsa_policy", "node_mode", "nodeMode", "reference_date", "today", "transit_date", "current_date", "entrypoint", "consult_entrypoint", ) _IDENTITY_KEYS = { "name", "email", "user_id", "userid", "user_email", "session_id", "sessionid", "full_name", "display_name", } def snapshot_identity_from_request(body: dict[str, Any] | None) -> dict[str, Any]: """Fields the BUG-727 key reads, copied without normalizing numbers. The consultation birth payload turns hour ``5`` into ``5.0``. Those are different cache keys, so the gate must look up the request's own values. """ payload = body if isinstance(body, dict) else {} return {key: payload[key] for key in _IDENTITY_FIELDS if key in payload} def negative_cache_dir() -> Path: raw = str(os.environ.get("JYOTISH_VEDASTRO_SNAPSHOT_NEGATIVE_CACHE_DIR") or "").strip() path = Path(raw) if raw else Path(__file__).resolve().parents[1] / "scratch" / "local" / "vedastro_snapshot_negative_cache" path.mkdir(parents=True, exist_ok=True) return path def _negative_path(cache_key: str) -> Path: if not re.fullmatch(r"[0-9a-f]{64}", cache_key): raise ValueError("vedastro negative cache key must be sha256 hex") return negative_cache_dir() / f"{cache_key}.json" def _strip_identity(value: Any) -> Any: if isinstance(value, dict): return { key: _strip_identity(item) for key, item in value.items() if str(key).strip().lower() not in _IDENTITY_KEYS } if isinstance(value, list): return [_strip_identity(item) for item in value] return value def _cache_body(body: dict[str, Any]) -> dict[str, Any]: identity = body.get("vedastro_snapshot_identity") if isinstance(identity, dict): return identity return body def _deferred_packet() -> dict[str, Any]: return { "scope": "vedastro_gateway_run", "status": "official_blocked", "official_closure_state": "official_blocked", "official_closure_reason": DEFERRED_IN_CONSULTATION, } def lookup_negative_snapshot( body: dict[str, Any] | None, *, today: str | None = None, ) -> dict[str, Any] | None: """Same-day negative hit only. Tomorrow is a different key, so it misses.""" payload = body if isinstance(body, dict) else {} served = (today or official_snapshot_reference_date(payload))[:10] path = _negative_path(snapshot_cache_key(payload, reference_date=served)) if not _SHA256_NAME.match(path.name) or not path.is_file(): return None try: record = json.loads(path.read_text(encoding="utf-8")) except (OSError, json.JSONDecodeError): return None if not isinstance(record, dict) or record.get("schema") != NEGATIVE_CACHE_SCHEMA: return None if str(record.get("reference_date") or "")[:10] != served: return None gateway = record.get("gateway") return dict(gateway) if isinstance(gateway, dict) else None def store_negative_snapshot( body: dict[str, Any] | None, gateway: dict[str, Any], *, today: str | None = None, ) -> dict[str, Any] | None: if not isinstance(gateway, dict): return None state = gateway.get("official_closure_state") or gateway.get("status") if state == "official_verified": return None payload = body if isinstance(body, dict) else {} served = (today or official_snapshot_reference_date(payload))[:10] try: datetime.strptime(served, "%Y-%m-%d") except ValueError: return None cache_key = snapshot_cache_key(payload, reference_date=served) record = { "schema": NEGATIVE_CACHE_SCHEMA, "cache_key": cache_key, "reference_date": served, "stored_at": datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ"), "gateway": _strip_identity(gateway), } path = _negative_path(cache_key) path.parent.mkdir(parents=True, exist_ok=True) encoded = json.dumps(record, ensure_ascii=False, sort_keys=True) fd, tmp_name = tempfile.mkstemp(prefix=f"{cache_key}.", suffix=".tmp", dir=str(path.parent)) try: with os.fdopen(fd, "w", encoding="utf-8") as handle: handle.write(encoded) os.replace(tmp_name, path) except Exception: with contextlib.suppress(OSError): os.unlink(tmp_name) raise return record def _verified_gateway(hit: dict[str, Any]) -> dict[str, Any] | None: record = hit.get("record") if isinstance(hit, dict) else None if not isinstance(record, dict): return None gateway = record.get("gateway") if not isinstance(gateway, dict): return None state = gateway.get("official_closure_state") or gateway.get("status") if state != "official_verified": return None if hit.get("freshness") == "stale": return annotate_stale_gateway( gateway, reference_date=str(record.get("reference_date") or ""), served_on_utc_date=str(hit.get("served_on_utc_date") or ""), ) return dict(gateway) def consult_rectification_vedastro_gateway(body: dict[str, Any] | None) -> dict[str, Any]: """Official snapshot for the consultation rectification gate. Never starts the runner.""" payload = body if isinstance(body, dict) else {} cache_body = _cache_body(payload) today = official_snapshot_reference_date(cache_body) reused = _verified_gateway(lookup_snapshot(cache_body, today=today) or {}) if reused is not None: return reused cached = lookup_negative_snapshot(cache_body, today=today) if cached is not None: return cached packet = _deferred_packet() with contextlib.suppress(OSError): store_negative_snapshot(cache_body, packet, today=today) return packet