#!/usr/bin/env python3
"""Fetch Boomi Data Integration data flows via API and save to a JSON file."""
from __future__ import annotations

import argparse
import datetime as dt
import json
import os
import sys
import urllib.parse
import urllib.request
from typing import Any

DEFAULT_BASE_URL = "https://api.rivery.io/v1"


def build_url(base_url: str, account_id: str, environment_id: str) -> str:
    safe_account_id = urllib.parse.quote(account_id, safe="")
    safe_environment_id = urllib.parse.quote(environment_id, safe="")
    return f"{base_url}/accounts/{safe_account_id}/environments/{safe_environment_id}/rivers"


def fetch_page(url: str, token: str) -> dict[str, Any]:
    request = urllib.request.Request(
        url,
        headers={
            "Authorization": f"Bearer {token}",
            "Accept": "application/json",
        },
    )
    try:
        with urllib.request.urlopen(  # nosec B310 - controlled endpoint
            request
        ) as response:
            return json.loads(response.read().decode("utf-8"))
    except urllib.error.HTTPError as exc:
        details = exc.read().decode("utf-8").strip()
        message = f"API request failed with HTTP {exc.code}."
        if details:
            message = f"{message} Response: {details}"
        raise RuntimeError(message) from exc


def fetch_all_pages(url: str, token: str) -> dict[str, Any]:
    combined: dict[str, Any] | None = None
    items: list[Any] = []
    next_url = url
    while next_url:
        payload = fetch_page(next_url, token)
        if combined is None:
            combined = dict(payload)
        page_items = payload.get("items", [])
        if not isinstance(page_items, list):
            raise ValueError("Unexpected API response: 'items' is not a list.")
        items.extend(page_items)
        next_url = payload.get("next_page")
    if combined is None:
        return {"items": []}
    combined["items"] = items
    combined["next_page"] = None
    combined["previous_page"] = None
    combined["page"] = 1
    combined["current_page_size"] = len(items)
    return combined


def parse_timestamp(value: Any) -> dt.datetime | None:
    if isinstance(value, (int, float)):
        return dt.datetime.fromtimestamp(value, tz=dt.timezone.utc)
    if isinstance(value, str):
        text = value.strip()
        if not text:
            return None
        if text.endswith("Z"):
            text = text[:-1] + "+00:00"
        try:
            parsed = dt.datetime.fromisoformat(text)
        except ValueError:
            return None
        if parsed.tzinfo is None:
            return parsed.replace(tzinfo=dt.timezone.utc)
        return parsed
    return None


def filter_recent_api_v2(payload: dict[str, Any]) -> dict[str, Any]:
    items = payload.get("items", [])
    if not isinstance(items, list):
        raise ValueError("Unexpected API response: 'items' is not a list.")
    cutoff = dt.datetime.now(dt.timezone.utc) - dt.timedelta(hours=24)
    filtered_items = []
    for item in items:
        if not isinstance(item, dict):
            continue
        if item.get("is_api_v2") is not True:
            continue
        pipelines = item.get("pipelines")
        has_recent_pipeline = False
        if isinstance(pipelines, list):
            for pipeline in pipelines:
                if not isinstance(pipeline, dict):
                    continue
                modified_at = parse_timestamp(pipeline.get("last_updated_at"))
                if modified_at and modified_at >= cutoff:
                    has_recent_pipeline = True
                    break
        if has_recent_pipeline:
            filtered_items.append(item)
            continue
        for key in ("last_updated_at", "created_at"):
            modified_at = parse_timestamp(item.get(key))
            if modified_at and modified_at >= cutoff:
                filtered_items.append(item)
                break
    filtered_payload = dict(payload)
    filtered_payload["items"] = filtered_items
    filtered_payload["next_page"] = None
    filtered_payload["previous_page"] = None
    filtered_payload["page"] = 1
    filtered_payload["current_page_size"] = len(filtered_items)
    return filtered_payload


def filter_api_v2(payload: dict[str, Any]) -> dict[str, Any]:
    items = payload.get("items", [])
    if not isinstance(items, list):
        raise ValueError("Unexpected API response: 'items' is not a list.")
    filtered_items = []
    for item in items:
        if not isinstance(item, dict):
            continue
        if item.get("is_api_v2") is not True:
            continue
        filtered_items.append(item)
    filtered_payload = dict(payload)
    filtered_payload["items"] = filtered_items
    filtered_payload["next_page"] = None
    filtered_payload["previous_page"] = None
    filtered_payload["page"] = 1
    filtered_payload["current_page_size"] = len(filtered_items)
    return filtered_payload


def write_output(path: str, payload: dict[str, Any]) -> None:
    output_dir = os.path.dirname(path)
    if output_dir:
        os.makedirs(output_dir, exist_ok=True)
    with open(path, "w", encoding="utf-8") as handle:
        json.dump(payload, handle, indent=2, ensure_ascii=False)
        handle.write("\n")


def parse_args(argv: list[str]) -> argparse.Namespace:
    parser = argparse.ArgumentParser(
        description="Fetch Boomi Data Integration data flows via API and save to a JSON file."
    )
    parser.add_argument(
        "--account-id",
        default=(
            os.getenv("BOOMI_ACCOUNT_ID")
            or os.getenv("BOOMI_ACCOUNT_ID_SECRET")
            or os.getenv("RIVERY_ACCOUNT_ID")
            or os.getenv("RIVERY_ACCOUNT_ID_SECRET")
        ),
        help=(
            "Boomi account ID (or set BOOMI_ACCOUNT_ID / BOOMI_ACCOUNT_ID_SECRET, "
            "or legacy RIVERY_ACCOUNT_ID / RIVERY_ACCOUNT_ID_SECRET)."
        ),
    )
    parser.add_argument(
        "--environment-id",
        default=(
            os.getenv("BOOMI_ENVIRONMENT_ID")
            or os.getenv("BOOMI_ENVIRONMENT_ID_SECRET")
            or os.getenv("RIVERY_ENVIRONMENT_ID")
            or os.getenv("RIVERY_ENVIRONMENT_ID_SECRET")
        ),
        help=(
            "Boomi environment ID (or set BOOMI_ENVIRONMENT_ID / "
            "BOOMI_ENVIRONMENT_ID_SECRET, or legacy RIVERY_ENVIRONMENT_ID / "
            "RIVERY_ENVIRONMENT_ID_SECRET)."
        ),
    )
    parser.add_argument(
        "--token",
        default=(
            os.getenv("BOOMI_API_TOKEN")
            or os.getenv("BOOMI_API_TOKEN_SECRET")
            or os.getenv("RIVERY_API_TOKEN")
            or os.getenv("RIVERY_API_TOKEN_SECRET")
        ),
        help=(
            "Boomi API token (or set BOOMI_API_TOKEN / BOOMI_API_TOKEN_SECRET, "
            "or legacy RIVERY_API_TOKEN / RIVERY_API_TOKEN_SECRET)."
        ),
    )
    parser.add_argument(
        "--base-url",
        default=DEFAULT_BASE_URL,
        help=f"API base URL (default: {DEFAULT_BASE_URL}).",
    )
    parser.add_argument(
        "--output",
        default=os.path.join("dataflows", "all-dataflows.json"),
        help=(
            "Path to the JSON file to create or update (default: "
            "dataflows/all-dataflows.json)."
        ),
    )
    parser.add_argument(
        "--changes-output",
        default=None,
        help=(
            "Path to the JSON file for filtered changes (default: "
            "<output-dir>/new_changes.json)."
        ),
    )
    parser.add_argument(
        "--api-v2-output",
        default=None,
        help=(
            "Path to the JSON file for API v2 data flows (default: "
            "<output-dir>/api_v2_dataflows.json)."
        ),
    )
    return parser.parse_args(argv)


def validate_args(args: argparse.Namespace) -> None:
    missing = []
    if not args.account_id:
        missing.append("--account-id or BOOMI_ACCOUNT_ID (or legacy RIVERY_ACCOUNT_ID)")
    if not args.environment_id:
        missing.append("--environment-id or BOOMI_ENVIRONMENT_ID (or legacy RIVERY_ENVIRONMENT_ID)")
    if not args.token:
        missing.append("--token or BOOMI_API_TOKEN (or legacy RIVERY_API_TOKEN)")
    if missing:
        raise ValueError(f"Missing required settings: {', '.join(missing)}")


def main(argv: list[str]) -> int:
    args = parse_args(argv)
    validate_args(args)

    url = build_url(args.base_url, args.account_id, args.environment_id)
    payload = fetch_all_pages(url, args.token)
    filtered_payload = filter_recent_api_v2(payload)
    api_v2_payload = filter_api_v2(payload)
    output_payload = {
        "source": "Boomi Data Integration API",
        "generated_at": dt.datetime.now(dt.timezone.utc).isoformat(),
        "account_id": args.account_id,
        "environment_id": args.environment_id,
        "data": payload,
    }
    output_changes_payload = {
        "source": "Boomi Data Integration API",
        "generated_at": dt.datetime.now(dt.timezone.utc).isoformat(),
        "account_id": args.account_id,
        "environment_id": args.environment_id,
        "filters": {
            "is_api_v2": True,
            "pipeline_last_updated_at_hours": 24,
        },
        "data": filtered_payload,
    }
    output_api_v2_payload = {
        "source": "Boomi Data Integration API",
        "generated_at": dt.datetime.now(dt.timezone.utc).isoformat(),
        "account_id": args.account_id,
        "environment_id": args.environment_id,
        "filters": {
            "is_api_v2": True,
        },
        "data": api_v2_payload,
    }
    write_output(args.output, output_payload)
    changes_output = args.changes_output
    if not changes_output:
        output_dir = os.path.dirname(args.output)
        changes_output = os.path.join(output_dir, "new_changes.json")
    write_output(changes_output, output_changes_payload)
    api_v2_output = args.api_v2_output
    if not api_v2_output:
        output_dir = os.path.dirname(args.output)
        api_v2_output = os.path.join(output_dir, "api_v2_dataflows.json")
    write_output(api_v2_output, output_api_v2_payload)
    print(f"Saved API response to {args.output}")
    print(f"Saved filtered changes to {changes_output}")
    print(f"Saved API v2 data flows to {api_v2_output}")
    return 0


if __name__ == "__main__":
    sys.exit(main(sys.argv[1:]))
