#!/usr/bin/env python3
"""Finframe JSON API: Python 3 stdlib only. No API key is written to disk.

FINFRAME_API_KEY='your private key' python3 finframe.py
Set the environment variable locally, never in a shared notebook or repository.
Successful requests consume quota; 304 responses do not. Default: one request.
"""
import gzip
import io
import json
import os
from pathlib import Path
import sys
from urllib.error import HTTPError
from urllib.parse import urlencode
from urllib.request import Request, urlopen

BASE = "https://fin.ksmzzang.duckdns.org"
IDS = "us_cpi,us_unemployment,us_treasury_10y"
STATE = Path("finframe-state.json")  # cursor, ETag and returned facts; never API keys
MAX_RESPONSE = 512 * 1024


class ApiRequestError(RuntimeError):
    def __init__(self, status, message):
        self.status = status
        super().__init__("HTTP %s: %s" % (status, message))


def get_json(path, key, etag=None):
    headers = {"Authorization": "Bearer " + key, "Accept": "application/json",
               "Accept-Encoding": "gzip"}
    if etag:
        headers["If-None-Match"] = etag
    try:
        with urlopen(Request(BASE + path, headers=headers), timeout=20) as response:
            raw = response.read(MAX_RESPONSE + 1)
            if len(raw) > MAX_RESPONSE:
                raise RuntimeError("API response exceeds the example's size limit.")
            encoding = response.headers.get("Content-Encoding", "").lower()
            if encoding == "gzip":
                with gzip.GzipFile(fileobj=io.BytesIO(raw)) as stream:
                    raw = stream.read(MAX_RESPONSE + 1)
            elif encoding not in ("", "identity"):
                raise RuntimeError("Unsupported response encoding.")
            if len(raw) > MAX_RESPONSE:
                raise RuntimeError("Decoded API response exceeds the example's size limit.")
            return json.loads(raw.decode("utf-8")), response.headers.get("ETag")
    except HTTPError as error:
        if error.code == 304:
            return None, etag
        try:
            detail = json.loads(error.read())
            message = detail.get("message") or detail.get("error", {})
            if isinstance(message, dict):
                message = message.get("message", message.get("code", ""))
        except (ValueError, OSError):
            message = "API request failed"
        # Do not print request headers or the key when reporting failures.
        raise ApiRequestError(error.code, message) from None


def pull_changes(state, key):
    """Apply selected series changes, at most 5 pages per run; resume next time.

    A separate history database should also apply changed_observations. This
    small example keeps the latest facts only, not a historical vintage archive.
    Changes do not contain collector health. Changed rows set freshness to None;
    run the default snapshot mode to obtain a fresh collector-health assessment.
    """
    selected = set(IDS.split(","))
    rows = {row["id"]: row for row in state["payload"]["data"]["series"]}
    changes_seen = 0
    for _ in range(5):
        query = urlencode({"cursor": state["cursor"], "limit": 100})
        payload, _ = get_json("/api/v1/changes?" + query, key)
        for change in payload["data"]["changes"]:
            if change["kind"] != "series" or change["id"] not in selected:
                continue
            if change["operation"] == "delete":
                rows.pop(change["id"], None)
            else:
                rows[change["id"]] = {**change["data"], "freshness": None}
            changes_seen += 1
        state["cursor"] = payload["data"]["next_cursor"]
        state["payload"]["meta"] = {**payload["meta"], "cursor": state["cursor"]}
        state["payload"]["data"]["series"] = list(rows.values())
        state["note"] = "Incremental fact cache; freshness requires a new snapshot."
        state.pop("etag", None)  # changes must not reuse an earlier snapshot ETag
        # Commit each completed page so the next run can safely resume.
        STATE.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8")
        if not payload["data"]["has_more"]:
            break
    print("Applied", changes_seen, "selected-series changes. Cursor saved.")


def main():
    key = os.environ.get("FINFRAME_API_KEY", "").strip()
    if not key:
        sys.exit("Set FINFRAME_API_KEY in your local environment first.")
    try:
        state = json.loads(STATE.read_text(encoding="utf-8")) if STATE.exists() else {}
    except (OSError, ValueError):
        state = {}
    if "--changes" in sys.argv and state.get("cursor") and state.get("payload"):
        try:
            pull_changes(state, key)
            return
        except ApiRequestError as error:
            if error.status != 410:
                raise
            print("Cursor expired; rebuilding from a fresh snapshot.")
            state = {}  # reset; never rebuild with an old ETag
    payload, etag = get_json("/api/v1/snapshot?ids=" + IDS, key, state.get("etag"))
    if payload is None:
        print("Unchanged (304). No request quota consumed.")
        return
    for row in payload["data"]["series"]:
        # Missing facts are null; do not silently convert them to zero.
        value = row.get("value")
        rendered = "unavailable" if value is None else str(value)
        print(row["id"], rendered, row["unit"]["label"], row.get("observation_period"))
        print("  Source:", row["source"]["url"])
        print("  Freshness:", row["freshness"]["status"])
    state = {"etag": etag, "cursor": payload["meta"].get("cursor"), "payload": payload}
    STATE.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8")
    print("Saved facts and cursor to", STATE)


if __name__ == "__main__":
    try:
        main()
    except (RuntimeError, OSError) as error:
        sys.exit(str(error))
