"""
Execução de jobs do AWS Glue no Data Lake.

Este módulo fornece funções para iniciar jobs do AWS Glue, 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 jobs sejam executados de forma confiável e consistente.

"""

from __future__ import annotations

import argparse
import json
import time
from pathlib import Path

from common import print_message
from src.datalake.commons import glue_utils
from src.datalake.commons.ci_utils import CYAN
from src.datalake.commons.data_env import load_runtime_env

WAIT_START_GLUE_JOB_WITH_PREDICATE_TIME = 30


def _prepare_glue_payload_full(
    use_predicate: str,
    truncate_data: str = "false",
    use_multithreading: str = "true",
    write_to_s3_enabled: str = "true",
) -> dict[str, str]:
    return {
        "--job-bookmark-option": "job-bookmark-disable",
        "--enable-metrics": "",
        "--enable-continuous-cloudwatch-log": "true",
        "--enable-job-insights": "true",
        "--enable-spark-ui": "true",
        "--enable-rename-algorithm-v2": "true",
        "--USE_PREDICATE": use_predicate,
        "--TRUNCATE_DATA": truncate_data,
        "--USE_MULTITHREADING": use_multithreading,
        "--WRITE_TO_S3_ENABLED": write_to_s3_enabled,
    }


def _write_payload(payload: dict[str, str]) -> None:
    Path("payload.json").write_text(
        json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8"
    )


def _load_env(phase: str, type_: str) -> dict[str, str]:
    env = load_runtime_env(apply_set_all=False, configure_aws=True)
    env["PHASE"] = phase
    if type_:
        env["SOURCE_SYSTEM"] = f"{env.get('SOURCE_SYSTEM', '')}_{type_}"
    return env


def _fix_jobs_list(jobs: list[str]) -> list[str]:
    jobs_list = []
    for job in jobs:
        job_fixed = job.strip()
        if job_fixed:
            if "," in job_fixed:
                job_list = [x.strip() for x in job_fixed.split(",") if x.strip()]
                jobs_list.extend(job_list)
            elif " " in job_fixed:
                job_list = [x.strip() for x in job_fixed.split(" ") if x.strip()]
                jobs_list.extend(job_list)
            else:
                jobs_list.append(job_fixed)
    return jobs_list


def _invoke_wait(
    env: dict[str, str],
    use_predicate: str,
    jobs: list[str],
    payload_builder,
    wait_between_runs: bool,
) -> None:
    print_message(f"[Dados] glue - invoke wait - jobs:  {', '.join(jobs)}")
    need_wait = False
    for use_predicate_fixed in [
        x.strip() for x in use_predicate.split(",") if x.strip()
    ]:
        if need_wait and wait_between_runs:
            print_message(
                "[Dados] glue - invoke wait - aguardando "
                f"{WAIT_START_GLUE_JOB_WITH_PREDICATE_TIME}s antes da próxima execução"
            )
            time.sleep(WAIT_START_GLUE_JOB_WITH_PREDICATE_TIME)

        payload = payload_builder(use_predicate_fixed)
        # _write_payload(payload)
        print_message("[Dados] glue - invoke wait - parâmetros")
        print_message(json.dumps(payload, ensure_ascii=False, indent=2))
        glue_utils.start_and_wait_glue_jobs(jobs, env, payload)
        need_wait = True


def _invoke_no_wait(
    env: dict[str, str], use_predicate: str, jobs: list[str], payload_builder
) -> None:
    print_message(f"[Dados] glue - invoke no wait - jobs:  {', '.join(jobs)}")
    for use_predicate_fixed in [
        x.strip() for x in use_predicate.split(",") if x.strip()
    ]:
        payload = payload_builder(use_predicate_fixed)
        # _write_payload(payload)
        print_message(
            f"[Dados] glue - invoke no wait - parâmetros \n {json.dumps(payload, ensure_ascii=False, indent=2)}"
        )
        for job in jobs:
            glue_utils.start_glue_job(job, env, payload)


def main() -> int:
    parser = argparse.ArgumentParser(description="Glue invoke")
    group = parser.add_mutually_exclusive_group(required=True)
    group.add_argument(
        "-w", "--wait", action="store_true", help="Wait for Glue jobs to complete"
    )
    group.add_argument(
        "-n",
        "--no-wait",
        action="store_true",
        help="Do not wait for Glue jobs to complete",
    )
    parser.add_argument(
        "-p",
        "--phase",
        required=True,
        choices=["transient", "raw", "stage", "analytics", "consolidate", "ml", "all"],
        help="Execution phase (transient, raw, stage, analytics, consolidate, ml, all)",
    )
    parser.add_argument(
        "-t",
        "--type",
        required=True,
        choices=["external", "anonymous", "tenant", "all"],
        help=(
            "Execution type (external, anonymous, tenant, or all) - "
            "will be used to set SOURCE_SYSTEM variable in the env"
        ),
    )
    parser.add_argument("--use_predicate", required=True)
    parser.add_argument(
        "--truncate_data",
        nargs="?",
        default="false",
        choices=["true", "false"],
        help="Whether to truncate data before processing (true or false)",
    )
    parser.add_argument(
        "--use_multithreading",
        nargs="?",
        default="true",
        choices=["true", "false"],
        help="Whether to use multithreading (true or false)",
    )
    parser.add_argument(
        "--write_to_s3_enabled",
        nargs="?",
        default="true",
        choices=["true", "false"],
        help="Whether to write to S3 (true or false)",
    )
    parser.add_argument("-j", "--jobs", nargs="*")

    args = parser.parse_args()

    print_message(f"{CYAN}[Dados] glue - invoke - args: {args}")
    try:
        env = _load_env(args.phase, args.type)
        if args.wait:
            print_message(
                f"{CYAN}[Dados] glue - invoke - espera acabar - usando parametros do job"
            )
            _invoke_wait(
                env,
                args.use_predicate,
                _fix_jobs_list(args.jobs),
                lambda predicate: _prepare_glue_payload_full(
                    predicate,
                    args.truncate_data or "false",
                    args.use_multithreading or "true",
                    args.write_to_s3_enabled or "true",
                ),
                wait_between_runs=True,
            )
        elif args.no_wait:
            print_message(
                f"{CYAN}[Dados] glue - invoke - não espera acabar - usando parametros passados"
            )
            _invoke_no_wait(
                env,
                args.use_predicate,
                _fix_jobs_list(args.jobs),
                lambda predicate: _prepare_glue_payload_full(
                    predicate,
                    args.truncate_data or "false",
                    args.use_multithreading or "true",
                    args.write_to_s3_enabled or "true",
                ),
            )
    except Exception as exc:
        print_message(str(exc))
        return 2
    return 0


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