"""
Executa ajustes de colunas em tabelas do Data Lake usando AWS Glue.

Este módulo fornece funções para iniciar jobs do AWS Glue que ajustam colunas em tabelas do Data Lake, aguardando a
conclusão dos jobs e gerenciando payloads de entrada e saída. As funções são projetadas para serem usadas em scripts de
CI/CD e automação de tarefas relacionadas ao Data Lake, garantindo que os ajustes de colunas sejam realizados de forma
confiável e consistente.

"""

from __future__ import annotations

import argparse
import json
import os
import time
from pathlib import Path

from common import print_message
from src.datalake.commons.aws_utils import (
    configure_aws_access,
    make_client,
    make_session,
)
from src.datalake.commons.data_env import export_to_process, load_data_ci_environment

WAIT_FOR_GLUE_JOB_TIME = 30


def _start_glue_job(
    env: dict[str, str], table_name: str, flow_type: str, phase: str, bucket_name: str
) -> str:
    job_name = f"{env.get('ENVIRONMENT')}_all_commons_adjust_columns"
    target_adjust_name = f"column_adjusts/{phase}/{table_name}.json"
    payload = {
        "--enable-metrics": "",
        "--enable-continuous-cloudwatch-log": "true",
        "--enable-rename-algorithm-v2": "true",
        "--USE_PREDICATE": "*",
        "--FLOW_TYPE": flow_type,
        "--MERGE_OLD_DATA": "true",
        "--TRUNCATE_DATA": "false",
        "--USE_MULTITHREADING": "false",
        "--WRITE_TO_S3_ENABLED": "false",
        "--PROCESS_ALL_AT_ONCE": "false",
        "--TARGET_BUCKET_NAME": bucket_name,
        "--TARGET_ADJUST_NAME": target_adjust_name,
    }

    Path("payload.json").write_text(
        json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8"
    )
    print_message(f"[Dados] glue start-job-run: {job_name}")

    session = make_session(env)
    glue = make_client(session, "glue", region_name=env.get("REGION", ""))
    response = glue.start_job_run(JobName=job_name, Arguments=payload)
    run_id = str(response.get("JobRunId", ""))
    if not run_id:
        raise RuntimeError("[Dados] Não foi possível obter JobRunId")
    return run_id


def _wait_glue_job(env: dict[str, str], run_id: str) -> None:
    job_name = f"{env.get('ENVIRONMENT')}_all_commons_adjust_columns"
    session = make_session(env)
    glue = make_client(session, "glue", region_name=env.get("REGION", ""))

    while True:
        response = glue.get_job_run(JobName=job_name, RunId=run_id)
        state = str(response.get("JobRun", {}).get("JobRunState", "")).upper()
        if state in {"ERROR", "TIMEOUT", "FAILED", "STOPPED"}:
            raise RuntimeError(
                f"[Dados] glue - Erro de execução encontrado, invalidando job! {state}"
            )
        if state in {"STARTING", "RUNNING", "STOPPING"}:
            print_message(
                f"[Dados] glue - {job_name} ainda em execução, tentando novamente em {WAIT_FOR_GLUE_JOB_TIME}s..."
            )
            time.sleep(WAIT_FOR_GLUE_JOB_TIME)
            continue
        print_message(f"[Dados] glue - {job_name} ok!")
        return


def main() -> int:
    parser = argparse.ArgumentParser(description="Adjust columns invoke")
    parser.add_argument("-a", action="store_true", dest="invoke")
    parser.add_argument("flow_type", nargs="?")
    parser.add_argument("phase", nargs="?")
    parser.add_argument("bucket_name", nargs="?")
    parser.add_argument("tables", nargs="*")
    args = parser.parse_args()

    if not args.invoke:
        parser.print_help()
        return 2

    try:
        env = load_data_ci_environment(os.environ.copy())
        env = configure_aws_access(env)

        if not args.flow_type or not args.phase or not args.bucket_name:
            raise RuntimeError(
                "[Dados] Parâmetros obrigatórios: FLOW_TYPE PHASE BUCKET_NAME"
            )

        bucket_name = args.bucket_name
        if env.get("ENVIRONMENT") != "prod" or env.get("DATAOFFICE_PROJECT_TYPE"):
            bucket_name = f"{bucket_name}-{env.get('ENVIRONMENT')}"

        export_to_process(env)
        for table_name in args.tables:
            run_id = _start_glue_job(
                env, table_name, args.flow_type, args.phase, bucket_name
            )
            _wait_glue_job(env, run_id)
    except Exception as exc:
        print_message(str(exc))
        return 2

    return 0


if __name__ == "__main__":
    raise SystemExit(main())
