import requests
from metaflow.exception import MetaflowException
from metaflow.metaflow_config import (
    ARGO_SERVER_USER_SA_SECRET_NAME,
    ARGO_WORKFLOWS_UI_URL,
    KUBERNETES_SERVICE_ACCOUNT,
)

from ..secret import get_secret


class ArgoClientException(MetaflowException):
    headline = "Argo Client error"


class ArgoClient:
    def __init__(self, namespace=None):
        self._namespace = namespace or "default"
        self._group = "argoproj.io"
        self._version = "v1alpha1"
        self.headers = self.create_headers()

    def create_headers(self):
        try:
            argo_server_user_sa_token = get_secret(ARGO_SERVER_USER_SA_SECRET_NAME)[
                "SA_TOKEN"
            ]
        except Exception as e:
            raise ArgoClientException(f"Error getting SA token from secret: {e}")
        headers = {"Authorization": f"Bearer {argo_server_user_sa_token}"}
        return headers

    def get_workflow(self, name):
        response = requests.get(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflows/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error getting workflow: {error_message}")
        return data

    def get_workflow_template(self, name):
        response = requests.get(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflow-templates/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(
                f"Error getting workflow template: {error_message}"
            )
        return data

    def get_workflow_templates(self):
        response = requests.get(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflow-templates/{self._namespace}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(
                f"Error getting workflow templates: {error_message}"
            )
        return data

    def register_workflow_template(self, name, workflow_template):
        # get existing workflow template
        existing_workflow_template = self.get_workflow_template(name)

        # set extra fields
        workflow_template["metadata"]["name"] = name
        workflow_template["spec"]["serviceAccountName"] = KUBERNETES_SERVICE_ACCOUNT

        # create if not found
        if existing_workflow_template is None:
            api_template = {
                "createOptions": {},
                "template": workflow_template,
            }
            response = requests.post(
                f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflow-templates/{self._namespace}",
                json=api_template,
                headers=self.headers,
            )

        # replace if found
        else:
            workflow_template["metadata"][
                "resourceVersion"
            ] = existing_workflow_template["metadata"]["resourceVersion"]
            api_template = {
                "createOptions": {},
                "template": workflow_template,
            }
            response = requests.put(
                f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflow-templates/{self._namespace}/{name}",
                json=api_template,
                headers=self.headers,
            )

        # check response
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(
                f"Error registering workflow template: {error_message}"
            )
        return data

    def delete_cronworkflow(self, name):
        """
        Issues an API call for deleting a cronworkflow

        Returns either the successful API response, or None in case the resource was not found.
        """
        response = requests.delete(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/cron-workflows/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error deleting cron workflow: {error_message}")
        return data

    def delete_workflow_template(self, name):
        """
        Issues an API call for deleting a cronworkflow

        Returns either the successful API response, or None in case the resource was not found.
        """
        response = requests.delete(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflow-templates/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(
                f"Error deleting workflow template: {error_message}"
            )
        return data

    def terminate_workflow(self, name):
        response = requests.put(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflows/{self._namespace}/{name}/terminate",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error terminating workflow: {error_message}")
        return data

    def suspend_workflow(self, name):
        response = requests.put(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflows/{self._namespace}/{name}/suspend",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error suspending workflow: {error_message}")
        return data

    def unsuspend_workflow(self, name):
        response = requests.put(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflows/{self._namespace}/{name}/resume",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error unsuspending workflow: {error_message}")
        return data

    def trigger_workflow_template(self, name, parameters={}):
        body = {
            "resourceKind": "WorkflowTemplate",
            "resourceName": name,
            "submitOptions": {
                "parameters": [f"{k}={v}" for k, v in parameters.items()]
            },
        }
        response = requests.post(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/workflows/{self._namespace}/submit",
            json=body,
            headers=self.headers,
        )
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error triggering workflow: {error_message}")
        return data

    def get_cronworkflow(self, name):
        response = requests.get(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/cron-workflows/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error getting cron workflow: {error_message}")
        return data

    def schedule_workflow_template(self, name, schedule=None, timezone=None):
        # get existing cronworkflow
        existing_cronworkflow = self.get_cronworkflow(name)

        # delete existing cronworkflow if found
        if existing_cronworkflow is not None:
            self.delete_cronworkflow(name)

        # skip if schedule is not provided
        if schedule is None:
            return None

        # create cronworkflow
        body = {
            "createOptions": {},
            "cronWorkflow": {
                "metadata": {"name": name},
                "spec": {
                    "suspend": schedule is None,
                    "schedule": schedule,
                    "timezone": timezone,
                    "workflowSpec": {"workflowTemplateRef": {"name": name}},
                },
            },
        }

        # create cronworkflow
        response = requests.post(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/cron-workflows/{self._namespace}",
            json=body,
            headers=self.headers,
        )
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error creating cron workflow: {error_message}")
        return data

    def get_sensor(self, name):
        response = requests.get(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/sensors/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error getting sensor: {error_message}")
        return data

    def register_sensor(self, name, sensor=None):
        if sensor is None:
            sensor = {}

        if not sensor:
            sensor["metadata"] = {}

        # get existing sensor
        existing_sensor = self.get_sensor(name)

        # create if not found
        if existing_sensor is None:
            api_template = {
                "createOptions": {},
                "sensor": dict(sensor),
            }
            response = requests.post(
                f"{ARGO_WORKFLOWS_UI_URL}/api/v1/sensors/{self._namespace}",
                json=api_template,
                headers=self.headers,
            )

        # replace if found
        else:
            api_template = {
                "createOptions": {},
                "sensor": dict(sensor),
            }
            api_template["sensor"]["metadata"]["resourceVersion"] = existing_sensor[
                "metadata"
            ]["resourceVersion"]
            response = requests.put(
                f"{ARGO_WORKFLOWS_UI_URL}/api/v1/sensors/{self._namespace}/{name}",
                json=api_template,
                headers=self.headers,
            )

        # check response
        data = response.json()
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error creating sensor: {error_message}")
        return data

    def delete_sensor(self, name):
        """
        Issues an API call for deleting a sensor

        Returns either the successful API response, or None in case the resource was not found.
        """
        # get existing sensor
        try:
            existing_sensor = self.get_sensor(name)
        except ArgoClientException as e:
            if "not found" not in str(e):
                raise e
            existing_sensor = None

        # skip if not found
        if existing_sensor is None:
            return None

        # delete if found
        response = requests.delete(
            f"{ARGO_WORKFLOWS_UI_URL}/api/v1/sensors/{self._namespace}/{name}",
            headers=self.headers,
        )
        data = response.json()
        if response.status_code == 404:
            return None
        if response.status_code != 200:
            error_message = data["message"]
            raise ArgoClientException(f"Error deleting sensor: {error_message}")
        return data
