#!/usr/bin/env python
"""Клиент публичного API rustata.ru.

Только стандартная библиотека; polars нужен единственной функции frame().
Токен берётся из окружения, в аргументы и в код он не попадает.
"""
from __future__ import annotations

import json
import os
import sys
import time
import urllib.error
import urllib.parse
import urllib.request
from pathlib import Path

BASE = os.environ.get("RUSTATA_API", "https://rustata.ru").rstrip("/")
MAX_SERIES_PAGE = 300      # больше одной страницей API не отдаёт (bulk_service.MAX_SERIES_LIMIT)
RETRY_ON_429 = 2           # ждём Retry-After и повторяем, но не бесконечно
MAX_PAGES = 20             # предохранитель: широкий фильтр иначе листает прод сотнями запросов

last_headers: dict[str, str] = {}   # заголовки последнего ответа в НИЖНЕМ регистре: квота, request-id


class TooManyPages(RuntimeError):
    """Фильтр слишком широкий. Это не ошибка сервера, а защита прода от сотен запросов."""


class ApiError(RuntimeError):
    """Ошибка в конверте API: code, message, request_id."""

    def __init__(self, status: int, code: str, message: str, request_id: str = ""):
        tail = f" (request_id={request_id})" if request_id else ""
        super().__init__(f"{status} {code}: {message}{tail}")
        self.status, self.code, self.message, self.request_id = status, code, message, request_id


def token() -> str:
    """RUSTATA_TOKEN из окружения, иначе строка в .env (он в .gitignore) или в ~/.rustata."""
    value = os.environ.get("RUSTATA_TOKEN", "").strip()
    if value:
        return value
    for env_file in (Path.cwd() / ".env", Path.home() / ".rustata"):
        if not env_file.exists():
            continue
        for line in env_file.read_text(encoding="utf-8").splitlines():
            name, _, raw = line.partition("=")
            if name.strip() == "RUSTATA_TOKEN" and raw.strip():
                return raw.strip().strip("'\"")
    raise SystemExit("Нет токена. Положите RUSTATA_TOKEN=rst_… в окружение или в .env; "
                     "выпустить - https://rustata.ru/account")


def request(path: str, params: dict | None = None, csv: bool = False) -> bytes:
    """GET к /api/v1. Возвращает тело ответа; ошибку разворачивает из конверта."""
    query = {k: v for k, v in (params or {}).items() if v not in (None, "")}
    url = f"{BASE}{path}" + (f"?{urllib.parse.urlencode(query)}" if query else "")
    headers = {"Authorization": f"Bearer {token()}",
               "Accept": "text/csv" if csv else "application/json",
               "User-Agent": "rustata-skill/1"}
    for attempt in range(RETRY_ON_429 + 1):
        try:
            with urllib.request.urlopen(urllib.request.Request(url, headers=headers), timeout=120) as response:
                last_headers.clear()
                # имена заголовков приводим к нижнему регистру: прод отдаёт
                # «X-Ratelimit-Remaining-Minute», а в документации написано «X-RateLimit-…»
                last_headers.update({k.lower(): v for k, v in response.headers.items()})
                return response.read()
        except urllib.error.HTTPError as exc:
            body = exc.read()
            try:
                envelope = json.loads(body)
                error = envelope.get("error") or {}
                code = error.get("code", "")
                message = error.get("message", "")
                req_id = envelope.get("request_id", "")
            except (ValueError, AttributeError):
                code, message, req_id = "", body[:300].decode("utf-8", "replace"), ""
            # 429 - минутный темп: ждём ровно столько, сколько сказано, и повторяем
            if exc.code == 429 and attempt < RETRY_ON_429:
                time.sleep(min(int(exc.headers.get("Retry-After") or 5), 120))
                continue
            raise ApiError(exc.code, code or f"HTTP_{exc.code}", message, req_id) from None
    raise ApiError(429, "RATE_LIMITED", "Минутный лимит не отпустил после повторов")


def get(path: str, **params) -> dict:
    return json.loads(request(path, params))


def limits() -> dict:
    """Что осталось по последнему ответу. У тарифов без месячной квоты значения None."""
    def number(name: str) -> int | None:
        raw = last_headers.get(name)
        return int(raw) if raw is not None else None

    return {"request_id": last_headers.get("x-request-id"),
            "values_limit": number("x-quota-values-limit"),
            "values_remaining": number("x-quota-values-remaining"),
            "minute_limit": number("x-ratelimit-limit-minute"),
            "minute_remaining": number("x-ratelimit-remaining-minute")}


def datasets() -> list[dict]:
    """Наборы: что подставлять в dataset, сколько рядов и за какой период."""
    return get("/api/v1/reference/datasets")["items"]


def geo(level: str = "") -> list[dict]:
    """Территории: код, название, canonical, уровень."""
    return get("/api/v1/reference/geo", level=level)["items"]


def find_geo(name: str, level: str = "region") -> list[tuple[str, str]]:
    """Коды территории по куску названия: (canonical, название).

    Пользоваться этим, а не догадкой по виду кода: canonical это «RU» плюс первые две цифры
    ОКАТО, и с ISO или автомобильным кодом он не совпадает.
    """
    needle = name.lower()
    return [(r["canonical"], r["name"]) for r in geo(level)
            if needle in str(r.get("name", "")).lower()]


def latest_period(dataset: str) -> str | None:
    """Последний период набора из справочника: дешевле, чем тянуть данные и смотреть максимум."""
    for row in datasets():
        if row["dataset"] == dataset or row["dataset"].endswith("/" + dataset):
            return row.get("latest_period")
    return None


def codes(dataset: str) -> list[dict]:
    """Коды продуктов ОДНОГО набора. Общего справочника нет: код 3203 в разных наборах разный."""
    return get("/api/v1/reference/codes", dataset=dataset)["items"]


def methodology(dataset: str = "") -> list[dict]:
    """Оговорки набора: охват, исключения, классификатор, пересматриваются ли значения."""
    return get("/api/v1/reference/methodology", dataset=dataset)["items"]


def updates(limit: int = 60) -> list[dict]:
    """Что и когда приезжало в склад: день, набор, добавленные периоды, пересмотренные точки."""
    return get("/api/v1/updates", limit=limit)["items"]


def search(q: str = "", limit: int = 20, **filters) -> list[dict]:
    """Поиск рядов индексом, с отбором по измерениям.

    filters: provider, dataset, geo, freq, class, transformation - через запятую «любое из».
    Строка поиска не обязательна: search(dataset="fedstat/indicator_31448", freq="M").
    """
    return get("/api/v1/series/search", q=q, limit=limit, **filters)["items"]


def observations(series_key: str, as_of: str = "latest", limit: int = 5000) -> list[dict]:
    return get(f"/api/v1/series/{series_key}/observations", as_of=as_of, limit=limit)["items"]


def heads(dataset: str) -> list[dict]:
    """Кандидаты в головной ряд набора: собирательные группировки и итоги.

    Агрегат в данных есть всегда, но называется не «итого»: у ИПП это «Собирательная
    классификационная группировка ... "Промышленность"». Суммировать разделы вместо него нельзя -
    у источника свой счёт, и наша сумма с ним не сойдётся.
    """
    marks = ("собирательн", "в целом", "всего", "итого", "промышленность")
    found = [c for c in codes(dataset)
             if any(m in str(c.get("name") or "").lower() for m in marks)]
    return sorted(found, key=lambda c: -(c.get("series_count") or 0))


def slice_data(dataset: str, max_pages: int = MAX_PAGES, active: bool = False, **params) -> dict:
    """Срез /api/v1/data со всеми страницами рядов сразу.

    Страница - это РЯДЫ, а не значения: каждый ряд приходит целиком по времени, склеивать
    историю не нужно. Потолок ответа - 60 000 значений на страницу, поэтому сужаем период
    или уменьшаем limit, а не ловим SLICE_TOO_BIG повтором.
    """
    params.setdefault("limit", MAX_SERIES_PAGE)
    offset, series, values, total = 0, [], [], 0
    status: dict = {}
    for _ in range(max_pages):
        page = get("/api/v1/data", dataset=dataset, offset=offset, **params)
        series += page["series"]
        values += page["observations"]
        status = page.get("dataset") or status      # статус набора один на все страницы
        total = page.get("series_total", len(series))
        if not page.get("has_more"):
            break
        offset = page["next_offset"]
    else:
        raise TooManyPages(
            f"Рядов под фильтр {total}, за {max_pages} страниц забрали {len(series)}. "
            f"Сузьте фильтр (geo, class, период) или поднимите max_pages осознанно: "
            f"каждая страница - отдельный запрос к проду.")
    if active:
        # Снятые с наблюдения в разбор трендов не берём: они тянут вниз долю растущих и
        # попадают в «без изменений». Отбор по признаку ряда, а не по наличию точек в окне.
        live = {s["series_key"] for s in series if s.get("status") != "stalled"}
        series = [s for s in series if s["series_key"] in live]
        values = [o for o in values if o["series_key"] in live]
    return {"series": series, "observations": values, "dataset": status,
            "series_count": len(series), "series_total": total,
            "observation_count": len(values)}


def frame(dataset: str, max_pages: int = MAX_PAGES, **params):
    """Тот же срез таблицей polars: ряды уже соединены со значениями.

    CSV несёт коды ВСЕХ измерений набора, поэтому таблица годится для join как есть
    (pandas в проекте не используем, ADR-0013).
    """
    import polars as pl

    params.setdefault("limit", MAX_SERIES_PAGE)
    frames, offset = [], 0
    for _ in range(max_pages):
        body = request("/api/v1/data", dict(params, dataset=dataset, offset=offset, format="csv"), csv=True)
        frames.append(pl.read_csv(body, separator=";", infer_schema_length=10_000))
        # у CSV нет полей страницы: идём, пока страница полная по рядам
        if frames[-1]["series_key"].n_unique() < int(params["limit"]):
            break
        offset += int(params["limit"])
    else:
        raise TooManyPages(f"Больше {max_pages} страниц. Сузьте фильтр (geo, class, период).")
    return pl.concat(frames, how="vertical_relaxed")


def _window_start(latest: str, freq: str, periods: int) -> str:
    """Начало окна на `periods` периодов назад от последнего. Без него срез тянет всю историю."""
    import datetime as dt

    try:
        last = dt.date.fromisoformat((latest + "-01")[:10] if len(latest) == 7 else latest)
    except ValueError:
        return ""
    step = {"W": 7, "D": 1}.get(freq)
    if step:
        return (last - dt.timedelta(days=step * periods)).isoformat()
    months = {"M": 1, "Q": 3, "S": 6, "H": 6, "A": 12, "Y": 12}.get(freq, 1) * periods
    year, month = divmod((last.year * 12 + last.month - 1) - months, 12)
    return f"{year:04d}-{month + 1:02d}-01"


def movers(dataset: str, geo: str = "643", periods: int = 12, start: str = "", end: str = "",
           top: int = 10, **filters) -> dict:
    """Сводка движения за последний период: кто вырос, кто упал, сколько рядов стоит.

    Один вызов вместо десятка: тянет срез, считает изменения и сам ограничивает окно последними
    `periods` периодами, иначе запрос уходит за всю историю набора.
    """
    import collections

    latest = latest_period(dataset) or ""
    freq = next((row.get("freq") or "" for row in datasets()
                 if row["dataset"] == dataset or row["dataset"].endswith("/" + dataset)), "")
    start = start or (_window_start(latest, str(freq), periods) if latest else "")

    filters.setdefault("class", "all")
    data = slice_data(dataset, geo=geo, start=start, end=end, labels=1, active=True, **filters)
    history: dict[str, dict[str, float]] = collections.defaultdict(dict)
    for o in data["observations"]:
        if o["value"] is not None:
            history[o["series_key"]][o["period"]] = float(o["value"])
    periods_seen = sorted({p for h in history.values() for p in h})
    if not periods_seen:
        return {"dataset": dataset, "periods": [], "items": [], "geo": geo}

    last = periods_seen[-1]
    prev = periods_seen[-2] if len(periods_seen) > 1 else None
    four = periods_seen[-5] if len(periods_seen) > 4 else periods_seen[0]
    multi_geo = len({s.get("geo") for s in data["series"]}) > 1

    # У наборов-индексов (ИПП, ИПЦ) значение УЖЕ процент: «рост на 141%» там получается из
    # 100.2 против 41.4, и это бессмыслица. Для них считаем разницу в пунктах, для уровней цен -
    # обычный прирост.
    index_like = {str(s.get("unit") or "").upper() for s in data["series"]} <= {"PCT", ""}

    def pct(now: float, was: float | None) -> float | None:
        if was in (None, 0):
            return None
        return (now - was) if index_like else (now / was - 1)

    items = []
    for s in data["series"]:
        points = history.get(s["series_key"], {})
        if last not in points:
            continue                      # ряда нет в последнем периоде: считать изменение не по чему
        name = s.get("class_label") or s.get("classification") or s.get("title")
        items.append({
            "name": f"{name} · {s.get('geo_label')}" if multi_geo else name,
            "code": s.get("classification"), "geo": s.get("geo_canonical"),
            "value": points[last],
            "change": pct(points[last], points.get(prev) if prev else None),
            "change_4": pct(points[last], points.get(four)),
            "change_window": pct(points[last], points.get(periods_seen[0])),
            "points": len(points),
            "flat": len(set(points.values())) == 1 and len(points) == len(periods_seen),
            # признак от API (см. A.3b); до выкладки его может не быть - тогда None
            "status": s.get("status"),
        })

    measures = {str(s.get("measure") or "") for s in data["series"]}
    moves = sorted([i["change"] for i in items if i["change"] is not None])
    middle = (moves[len(moves) // 2] if len(moves) % 2 else
              (moves[len(moves) // 2 - 1] + moves[len(moves) // 2]) / 2) if moves else None
    ranked = sorted([i for i in items if i["change"] is not None], key=lambda i: i["change"])
    return {
        "dataset": dataset, "geo": geo, "periods": periods_seen, "last": last, "prev": prev,
        "items": items, "count": len(items),
        "up": sum(1 for m in moves if m > 0), "down": sum(1 for m in moves if m < 0),
        "same": sum(1 for m in moves if m == 0),
        "median": middle, "mean": (sum(moves) / len(moves)) if moves else None,
        "index_like": index_like, "measures": sorted(m for m in measures if m),
        "flat": [i["name"] for i in items if i["flat"]],
        "stalled": [i["name"] for i in items if i.get("status") == "stalled"],
        "dataset_status": (data.get("dataset") or {}).get("status"),
        "top_up": ranked[-top:][::-1], "top_down": ranked[:top],
    }


def regions(dataset: str, code: str, level: str = "region", periods: int = 12,
            start: str = "", end: str = "", top: int = 10, **filters) -> dict:
    """Один продукт по территориям: кто дорожает, кто дешевеет, каков разброс уровня.

    Второй типовой вопрос после «что интересного за неделю», и вручную он стоит десятка запусков:
    вытащить срез, соединить ряды со значениями, посчитать изменения, отсортировать.
    """
    import collections

    latest = latest_period(dataset) or ""
    freq = next((row.get("freq") or "" for row in datasets()
                 if row["dataset"] == dataset or row["dataset"].endswith("/" + dataset)), "")
    start = start or (_window_start(latest, str(freq), periods) if latest else "")

    data = slice_data(dataset, geo=level, start=start, end=end, labels=1, active=True,
                      **dict(filters, **{"class": code}))
    history: dict[str, dict[str, float]] = collections.defaultdict(dict)
    for o in data["observations"]:
        if o["value"] is not None:
            history[o["series_key"]][o["period"]] = float(o["value"])
    periods_seen = sorted({p for h in history.values() for p in h})
    if not periods_seen:
        return {"dataset": dataset, "code": code, "items": [], "periods": []}

    last, first = periods_seen[-1], periods_seen[0]
    prev = periods_seen[-2] if len(periods_seen) > 1 else None
    items = []
    for s in data["series"]:
        points = history.get(s["series_key"], {})
        if last not in points:
            continue
        items.append({
            "geo": s.get("geo_canonical"), "name": s.get("geo_label"),
            "value": points[last],
            "change": (points[last] / points[prev] - 1) if prev and points.get(prev) else None,
            "change_window": (points[last] / points[first] - 1) if points.get(first) else None,
            "status": s.get("status"),
            "label": s.get("class_label") or code,
        })
    ranked = sorted([i for i in items if i["change_window"] is not None],
                    key=lambda i: i["change_window"])
    levels = sorted(i["value"] for i in items)
    return {"dataset": dataset, "code": code, "label": items[0]["label"] if items else code,
            "periods": periods_seen, "last": last, "items": items, "count": len(items),
            "cheapest": ranked and min(items, key=lambda i: i["value"]),
            "dearest": ranked and max(items, key=lambda i: i["value"]),
            "median_level": levels[len(levels) // 2] if levels else None,
            "top_up": ranked[-top:][::-1], "top_down": ranked[:top]}


def render_regions(summary: dict) -> str:
    if not summary.get("items"):
        return f"{summary['dataset']} / {summary['code']}: рядов не нашлось"
    pct = lambda v: "     -" if v is None else f"{v * 100:+6.2f}%"   # noqa: E731
    cheap, dear = summary["cheapest"], summary["dearest"]
    lines = [
        f"{summary['label']} ({summary['code']}) | {summary['dataset']} | территорий "
        f"{summary['count']} | {summary['periods'][0]}..{summary['last']}",
        f"уровень: медиана {summary['median_level']:,.2f}".replace(",", " ")
        + f" | дешевле всего {cheap['name']} {cheap['value']:,.2f}".replace(",", " ")
        + f" | дороже всего {dear['name']} {dear['value']:,.2f}".replace(",", " "),
        "",
        f"{'период':>8} {'окно':>8}  значение     территория",
    ]
    for title, rows in (("рост", summary["top_up"]), ("падение", summary["top_down"])):
        lines.append(f"-- {title} за окно --")
        for i in rows:
            lines.append(f"{pct(i['change'])} {pct(i['change_window']):>8}  {i['value']:>9,.2f}"
                         f"  {i['name']}".replace(",", " "))
    return chr(10).join(lines)


def render_movers(summary: dict) -> str:
    """Сводка текстом - компактно, чтобы читалась целиком без дополнительных запросов."""
    if not summary.get("items"):
        return f"{summary['dataset']}: под фильтры ничего не попало"
    index_like = summary.get("index_like")
    # У индексов печатаем пункты: «+1.40 п.п.» честнее, чем «+1.40%» от величины, которая сама процент
    pct = ((lambda v: "        -" if v is None else f"{v:+6.2f} п.п.") if index_like
           else (lambda v: "     -" if v is None else f"{v * 100:+6.2f}%"))     # noqa: E731
    warn = []
    if index_like:
        warn.append("значения - индексы (уже проценты), изменения показаны в пунктах")
    if len(summary.get("measures") or []) > 1:
        warn.append("в наборе несколько показателей ({}), сравнивать их между собой нельзя - "
                    "сузьте measure=".format(", ".join(summary["measures"])))
    lines = [
        f"{summary['dataset']} | территория {summary['geo']} | периодов {len(summary['periods'])}"
        f" ({summary['periods'][0]}..{summary['last']})",
        f"рядов {summary['count']} | подорожало {summary['up']}, подешевело {summary['down']}, "
        f"без изменений {summary['same']} | медиана {pct(summary['median'])}, "
        f"среднее {pct(summary['mean'])}",
        # «снят с наблюдения» и «стоит на месте» - разные вещи: первое говорит API по всей
        # истории набора, второе видно в окне и часто настоящее (регулируемые тарифы).
        f"набор: {summary.get('dataset_status') or 'статус не пришёл'}"
        + (f" | снято с наблюдения: {len(summary['stalled'])}" if summary.get("stalled") else ""),
        f"стоят весь период: {len(summary['flat'])}"
        + (f" ({', '.join(summary['flat'][:5])}{'…' if len(summary['flat']) > 5 else ''})"
           if summary["flat"] else ""),
        *[f"ВНИМАНИЕ: {w}" for w in warn],
        "",
        f"{'период':>8} {'4 периода':>10} {'окно':>8}  значение     ряд",
    ]
    for title, rows in (("рост", summary["top_up"]), ("падение", summary["top_down"])):
        lines.append(f"-- {title} --")
        for i in rows:
            lines.append(f"{pct(i['change'])} {pct(i['change_4']):>10} {pct(i['change_window']):>8}"
                         f"  {i['value']:>9,.2f}  {i['name']}".replace(",", " "))
    return chr(10).join(lines)


_USAGE = """Использование:
  rustata.py datasets
  rustata.py geo [level=region]
  rustata.py codes <dataset>
  rustata.py heads <dataset>                                      головной ряд набора (агрегат)
  rustata.py methodology [<dataset>]                              оговорки набора и пересмотры
  rustata.py updates [limit=20]                                   что нового приехало
  rustata.py search <строка> [provider=fedstat] [dataset=…] [geo=RU-45] [freq=M]
  rustata.py data <dataset> [geo=region] [class=3203] [start=2024-01] [end=…] [labels=1]
  rustata.py movers <dataset> [geo=RU-45] [periods=12] [top=10] [measure=YOY]  сводка движения
  rustata.py regions <dataset> <код продукта> [level=region] [periods=12]  один продукт по территориям
  rustata.py get <путь> [параметр=значение …]
"""


def _main(argv: list[str]) -> int:
    if not argv or argv[0] in ("-h", "--help"):
        print(_USAGE)
        return 0
    command, rest = argv[0], argv[1:]
    named = dict(a.split("=", 1) for a in rest if "=" in a)
    plain = [a for a in rest if "=" not in a]
    if command == "datasets":
        result = datasets()
    elif command == "geo":
        result = geo(named.get("level", ""))
    elif command == "codes":
        result = codes(plain[0])
    elif command == "heads":
        result = heads(plain[0])
    elif command == "methodology":
        result = methodology(plain[0] if plain else "")
    elif command == "updates":
        result = updates(int(named.get("limit", 60)))
    elif command == "search":
        result = search(" ".join(plain), int(named.pop("limit", 20)), **named)
    elif command == "data":
        result = slice_data(plain[0], **named)
    elif command == "regions":
        # Остальные name=value уходят фильтрами измерений (measure=, unit=, transformation=):
        # молча их терять нельзя, иначе в рейтинге смешиваются несравнимые ряды.
        level = named.pop("level", "region")
        print(render_regions(regions(plain[0], plain[1], level=level,
                                     periods=int(named.pop("periods", 12)),
                                     start=named.pop("start", ""), end=named.pop("end", ""),
                                     top=int(named.pop("top", 10)), **named)))
        return 0
    elif command == "movers":
        summary = movers(plain[0], geo=named.pop("geo", "643"),
                         periods=int(named.pop("periods", 12)), start=named.pop("start", ""),
                         end=named.pop("end", ""), top=int(named.pop("top", 10)), **named)
        print(render_movers(summary))
        return 0
    elif command == "get":
        result = get(plain[0], **named)
    else:
        print(_USAGE)
        return 2
    print(json.dumps(result, ensure_ascii=False, indent=2, default=str))
    left = limits()
    print(f"# осталось: значений {left['values_remaining'] or 'без квоты'}, "
          f"запросов в минуту {left['minute_remaining']}; {left['request_id']}", file=sys.stderr)
    return 0


if __name__ == "__main__":
    # Консоль Windows по умолчанию cp1252: без этого кириллица в выводе роняет запуск.
    for _stream in (sys.stdout, sys.stderr):
        try:
            _stream.reconfigure(encoding="utf-8")
        except AttributeError:
            pass
    try:
        sys.exit(_main(sys.argv[1:]))
    except ApiError as error:
        print(error, file=sys.stderr)
        sys.exit(1)
