| """Live-demo ASGI entrypoint for the self-reflective-apis research server. |
| |
| uvicorn live_app:app --host 127.0.0.1 --port 8100 |
| |
| This file is OURS. Nothing inside ``_clone/`` is modified, ever -- it is only |
| imported from, and extra routes are composed onto the upstream FastAPI app |
| object at runtime. |
| |
| What it does |
| ------------ |
| 1. Puts ``_clone/apis`` on ``sys.path`` and imports the upstream recipe app |
| (``recipe.main:app``), the self-contained domain on the pinned default |
| branch. |
| 2. Imports the Acme ``required_promo`` policy module, which registers the Acme |
| CSM/promo table ``active_csm_codes`` in the shared ``TableRegistry`` |
| singleton. Both apps and every ``table_versions`` block read that same |
| singleton, so the Acme table becomes visible in recipe responses and |
| reloadable through the recipe app's ``/admin/reload-table`` -- without |
| importing the Acme FastAPI app, which does not import at this commit |
| (see RUNBOOK_live.md, "Upstream surprises"). |
| 3. Rebinds the registry path for the drifting table from the read-only clone |
| to a writable copy under ``runtime_data/`` so that drift rotation never |
| writes into ``_clone/``. ``drift_cron.py`` swaps that same file. Which |
| table drifts is chosen by the LIVE_DRIFT_* variables; the default is the |
| recipe table the API reads when it builds a recovery suggestion, so a |
| rotation can change an answer a client is already holding. |
| 4. Adds two routes of our own: ``GET /live/status``, describing what is loaded |
| and on what public cadence it drifts, and ``GET /live``, a read-only status |
| page that renders the same payload. |
| |
| The live table file is the single source of truth for drift position: its |
| ``version`` major number is the snapshot index. There is no separate cursor to |
| get out of sync with, and a restart never rewinds a rotation. |
| """ |
| from __future__ import annotations |
|
|
| import json |
| import os |
| import shutil |
| import sys |
| from datetime import datetime, timezone |
| from pathlib import Path |
|
|
| LIVE_ROOT = Path(__file__).resolve().parent |
|
|
| |
| |
| |
| |
| |
| sys.dont_write_bytecode = True |
| sys.pycache_prefix = str(LIVE_ROOT / "runtime_data" / "pycache") |
| os.environ.setdefault("PYTHONDONTWRITEBYTECODE", "1") |
|
|
| |
| |
| for _stream in (sys.stdout, sys.stderr): |
| try: |
| _stream.reconfigure(encoding="utf-8", errors="replace") |
| except (AttributeError, ValueError): |
| pass |
|
|
| UPSTREAM_ROOT = Path(os.environ.get("LIVE_UPSTREAM_ROOT", LIVE_ROOT / "_clone")).resolve() |
| CLONE_APIS = UPSTREAM_ROOT / "apis" |
| SNAPSHOT_DIR = Path(os.environ.get("LIVE_SNAPSHOT_DIR", LIVE_ROOT / "snapshots")).resolve() |
| RUNTIME_DATA = Path(os.environ.get("LIVE_RUNTIME_DATA", LIVE_ROOT / "runtime_data")).resolve() |
| |
| |
| RESET_TABLES = os.environ.get("LIVE_RESET_TABLES", "") == "1" |
|
|
| |
| |
| |
| |
| TABLE_NAME = os.environ.get("LIVE_DRIFT_TABLE", "incompatible_ingredients") |
| FILE_STEM = os.environ.get("LIVE_DRIFT_FILE", "incompatible_combinations") |
| MIRROR_DIR = os.environ.get("LIVE_DRIFT_MIRROR", "recipe") |
| SRC_REL = os.environ.get("LIVE_DRIFT_SRC", "recipe/data") |
| |
| SEED_INDEX = int(os.environ.get("LIVE_SEED_INDEX", "1")) |
| LIVE_TABLE = RUNTIME_DATA / MIRROR_DIR / (FILE_STEM + ".json") |
|
|
| |
| |
| DRIFT_SCHEDULE = os.environ.get("LIVE_DRIFT_SCHEDULE", "daily at 03:00 host local time") |
|
|
| |
| _UPSTREAM_DATA = CLONE_APIS.joinpath(*SRC_REL.split("/")) |
| _UPSTREAM_SNAPSHOTS = _UPSTREAM_DATA / "snapshots" |
|
|
| if not CLONE_APIS.is_dir(): |
| raise SystemExit( |
| "Upstream clone not found at {}. Run:\n" |
| " git clone https://github.com/arquicanedo/self-reflective-apis _clone".format(CLONE_APIS) |
| ) |
| if str(CLONE_APIS) not in sys.path: |
| sys.path.insert(0, str(CLONE_APIS)) |
|
|
|
|
| def seed_snapshots() -> list[Path]: |
| """Mirror the upstream snapshot ring into snapshots/ on first run. |
| |
| One-way copy (clone -> ours); nothing is ever written back. If snapshots/ |
| is already populated, it is left alone so an operator can extend or edit |
| the ring without the server clobbering it. |
| """ |
| SNAPSHOT_DIR.mkdir(parents=True, exist_ok=True) |
| ring = sorted(SNAPSHOT_DIR.glob(FILE_STEM + "_v*.json")) |
| if ring: |
| return ring |
|
|
| shipped = _UPSTREAM_DATA / (FILE_STEM + ".json") |
| if shipped.is_file(): |
| |
| shutil.copyfile(shipped, SNAPSHOT_DIR / (FILE_STEM + "_v1.json")) |
| for src in sorted(_UPSTREAM_SNAPSHOTS.glob(FILE_STEM + "_v*.json")): |
| shutil.copyfile(src, SNAPSHOT_DIR / src.name) |
|
|
| ring = sorted(SNAPSHOT_DIR.glob(FILE_STEM + "_v*.json")) |
| if len(ring) < 2: |
| raise SystemExit( |
| "Need at least two snapshots of {!r} in {} for drift to mean anything.".format( |
| TABLE_NAME, SNAPSHOT_DIR |
| ) |
| ) |
| return ring |
|
|
|
|
| def seed_live_table() -> Path: |
| """Create the writable live table if absent. |
| |
| An existing file is never overwritten unless LIVE_RESET_TABLES=1, so a |
| server restart resumes at whatever snapshot drift_cron last rotated to. |
| """ |
| LIVE_TABLE.parent.mkdir(parents=True, exist_ok=True) |
| if RESET_TABLES or not LIVE_TABLE.exists(): |
| shutil.copyfile(SNAPSHOT_DIR / (FILE_STEM + "_v%d.json" % SEED_INDEX), LIVE_TABLE) |
| return LIVE_TABLE |
|
|
|
|
| |
| |
| |
| _ring = seed_snapshots() |
| seed_live_table() |
|
|
| from recipe.main import app |
| from shared.table_registry import registry |
|
|
| try: |
| from acme_billing.policies import required_promo |
| ACME_TABLE_LOADED = True |
| except Exception as _exc: |
| ACME_TABLE_LOADED = False |
| print("[live] WARNING: Acme promo table not loaded: %r" % (_exc,), file=sys.stderr) |
|
|
| |
| |
| if TABLE_NAME in registry._paths: |
| registry._paths[TABLE_NAME] = LIVE_TABLE |
| registry.reload(TABLE_NAME) |
| else: |
| registry.load(TABLE_NAME, LIVE_TABLE) |
|
|
| print( |
| "[live] %s v%s from %s | ring of %d | drift: %s" |
| % (TABLE_NAME, registry.version(TABLE_NAME), LIVE_TABLE, len(_ring), DRIFT_SCHEDULE), |
| file=sys.stderr, |
| flush=True, |
| ) |
|
|
|
|
| from fastapi.responses import HTMLResponse |
|
|
| _STARTED = datetime.now(timezone.utc) |
| _ROTATIONS = [] |
| _LAST = {"version": None, "data": None} |
| _PAGE = LIVE_ROOT / "live_page.html" |
| |
| |
| |
| _VISITOR_CALLS = {"n": 0} |
|
|
|
|
| @app.middleware("http") |
| async def _count_visitor_calls(request, call_next): |
| if request.url.path.startswith("/api/") and "x-live-probe" not in request.headers: |
| _VISITOR_CALLS["n"] += 1 |
| return await call_next(request) |
|
|
|
|
| def _iso(dt): |
| return dt.astimezone(timezone.utc).isoformat(timespec="seconds").replace("+00:00", "Z") |
|
|
|
|
| def _upstream_commit(): |
| """Short-circuit git: read .git/HEAD by hand rather than shelling out.""" |
| git = UPSTREAM_ROOT / ".git" |
| try: |
| head = (git / "HEAD").read_text(encoding="utf-8").strip() |
| except OSError: |
| return None |
| if not head.startswith("ref: "): |
| return head |
| ref = head[5:].strip() |
| try: |
| return (git / ref).read_text(encoding="utf-8").strip() |
| except OSError: |
| pass |
| try: |
| for line in (git / "packed-refs").read_text(encoding="utf-8").splitlines(): |
| if line.endswith(" " + ref): |
| return line.split(" ", 1)[0] |
| except OSError: |
| pass |
| return None |
|
|
|
|
| UPSTREAM_COMMIT = _upstream_commit() |
|
|
|
|
| def _read_live_table(): |
| try: |
| blob = json.loads(LIVE_TABLE.read_text(encoding="utf-8")) |
| except (OSError, json.JSONDecodeError): |
| return None, {} |
| return blob.get("version"), blob.get("data") or {} |
|
|
|
|
| def _diff_entries(before, after): |
| """Diff two table payloads. |
| |
| Scalar columns report old/new. Row lists report added/removed, matching the |
| shape the upstream reload endpoint returns for the same tables. |
| """ |
| changed = {} |
| for key in set(before) | set(after): |
| old, new = before.get(key), after.get(key) |
| if old == new: |
| continue |
| if isinstance(old, list) or isinstance(new, list): |
| old_rows, new_rows = old or [], new or [] |
| changed[key] = { |
| "removed": [r for r in old_rows if r not in new_rows], |
| "added": [r for r in new_rows if r not in old_rows], |
| } |
| else: |
| changed[key] = {"old": old, "new": new} |
| return changed |
|
|
|
|
| def _track(served, on_disk, entries): |
| """Record a rotation once the reload has caught up with the file. |
| |
| While the file is ahead of the served copy the cached 'before' state is |
| left alone, so the diff is taken against what was actually being served. |
| """ |
| if on_disk != served: |
| return |
| previous = _LAST["version"] |
| if previous is not None and served != previous: |
| before = _LAST["data"] or {} |
| changed = _diff_entries(before, entries) |
| try: |
| when = datetime.fromtimestamp(LIVE_TABLE.stat().st_mtime, timezone.utc) |
| except OSError: |
| when = datetime.now(timezone.utc) |
| _ROTATIONS.insert(0, {"ts": _iso(when), "table": TABLE_NAME, |
| "from": previous, "to": served, "changed": changed}) |
| del _ROTATIONS[12:] |
| _LAST["version"] = served |
| _LAST["data"] = entries |
|
|
|
|
| _seed_version, _seed_entries = _read_live_table() |
| _track(registry.version(TABLE_NAME), _seed_version, _seed_entries) |
|
|
|
|
| @app.get("/live/status") |
| def live_status(): |
| """Public status: what is loaded, what drifts, and on what cadence.""" |
| served = registry.version(TABLE_NAME) |
| on_disk, entries = _read_live_table() |
| _track(served, on_disk, entries) |
| return { |
| "service": "self-reflective-apis live demo", |
| "upstream": "github.com/arquicanedo/self-reflective-apis (default branch, pinned)", |
| "upstream_commit": UPSTREAM_COMMIT, |
| "started_at": _iso(_STARTED), |
| "table_versions": registry.all_versions(), |
| "drift": { |
| "table": TABLE_NAME, |
| "schedule": DRIFT_SCHEDULE, |
| "ring": [p.name for p in sorted(SNAPSHOT_DIR.glob(FILE_STEM + "_v*.json"))], |
| "served_version": served, |
| "on_disk_version": on_disk, |
| "reload_pending": on_disk is not None and on_disk != served, |
| "entries": entries, |
| }, |
| "rotations": _ROTATIONS, |
| "visitor_requests": _VISITOR_CALLS["n"], |
| "acme_promo_table_loaded": ACME_TABLE_LOADED, |
| "admin": "routes under /admin require a bearer token at the gate (port 8180)", |
| } |
|
|
|
|
| @app.get("/live", include_in_schema=False) |
| @app.get("/live/", include_in_schema=False) |
| def live_page(): |
| return HTMLResponse(_PAGE.read_text(encoding="utf-8"), |
| headers={"Cache-Control": "no-store"}) |
|
|