from __future__ import annotations

import os
from datetime import datetime
from threading import Thread
from typing import Any

from app.core.logging import (
    bind_log_context,
    clear_log_context,
    get_audit_logger,
    get_log_context,
    set_log_context,
)
from app.rda.application.fhir_validation import ValidateRdaFhirProjectionUseCase
from app.rda.application.models import GenerateRdaCommand, QueuedRdaGeneration, RdaStatusResult
from app.rda.application.models import RdaValidationReviewResult
from app.rda.application.use_cases import GenerateCaseRdaUseCase
from app.rda.domain.models import RdaArtifactType, RdaJobState
from app.rda.infrastructure.providers import EpicrisisServiceCaseContextProvider
from app.rda.infrastructure.repositories import MongoRdaArtifactRepository, MongoRdaJobStatusRepository


class RdaService:
    """Fachada para RDA interno versionado y asíncrono."""

    def __init__(
        self,
        *,
        mongo_analyses: Any,
        colombia_tz: Any,
        case_epicrisis_service: Any,
    ):
        self.audit_logger = get_audit_logger()
        self.colombia_tz = colombia_tz
        self.artifact_repository = MongoRdaArtifactRepository(mongo_analyses, colombia_tz)
        self.job_status_repository = MongoRdaJobStatusRepository(mongo_analyses, colombia_tz)
        self.context_provider = EpicrisisServiceCaseContextProvider(case_epicrisis_service)
        self.generate_case_rda_use_case = GenerateCaseRdaUseCase(
            context_provider=self.context_provider,
            artifact_repository=self.artifact_repository,
            job_status_repository=self.job_status_repository,
        )
        self.validate_rda_fhir_projection_use_case = ValidateRdaFhirProjectionUseCase()

    def ensure_indexes(self) -> None:
        self.artifact_repository.ensure_indexes()
        self.job_status_repository.ensure_indexes()

    def _normalize_artifact_type(self, artifact_type: str | RdaArtifactType) -> RdaArtifactType:
        if isinstance(artifact_type, RdaArtifactType):
            return artifact_type
        return RdaArtifactType(str(artifact_type).strip().lower())

    def _job_state_value(self, status: str | RdaJobState) -> str:
        return status.value if isinstance(status, RdaJobState) else str(status)

    def _queue_status(
        self,
        *,
        username: str,
        case_key: str,
        artifact_type: RdaArtifactType,
        job_id: str,
        reused: bool,
    ) -> QueuedRdaGeneration:
        self.job_status_repository.upsert_status(
            username=username,
            case_key=case_key,
            artifact_type=artifact_type,
            status=RdaJobState.QUEUED.value,
            job_id=job_id,
            reused=reused,
        )
        return QueuedRdaGeneration(
            case_key=case_key,
            artifact_type=artifact_type,
            job_id=job_id,
            status=RdaJobState.QUEUED,
            reused=reused,
        )

    def _spawn_inprocess_job(
        self,
        *,
        username: str,
        case_key: str,
        artifact_type: RdaArtifactType,
        force: bool,
        job_id: str,
    ) -> None:
        audit_context = get_log_context()

        def _runner() -> None:
            self.run_generation_job(
                username,
                case_key,
                artifact_type,
                force=force,
                job_id=job_id,
                audit_context=audit_context,
            )

        Thread(target=_runner, daemon=True).start()

    def queue_case_rda(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
        *,
        force: bool = False,
    ) -> QueuedRdaGeneration:
        normalized_type = self._normalize_artifact_type(artifact_type)
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        current_status = self.get_case_rda_status(username, case_key, normalized_type)
        if current_status and current_status.status in {RdaJobState.QUEUED, RdaJobState.PROCESSING}:
            reused_result = QueuedRdaGeneration(
                case_key=case_key,
                artifact_type=normalized_type,
                job_id=str(current_status.job_id or ""),
                status=current_status.status,
                reused=True,
            )
            self.audit_logger.business_event(
                event_type="rda.queue_reused",
                action="queue_case_rda",
                outcome="success",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                    "job_id": reused_result.job_id,
                    "status": self._job_state_value(reused_result.status),
                    "reused": reused_result.reused,
                },
            )
            return reused_result

        dispatcher_mode = os.getenv("BATCH_DISPATCHER", "inprocess").strip().lower()
        if dispatcher_mode == "celery":
            from app.batch_processing.celery_app import generate_rda_job

            result = generate_rda_job.delay(
                username,
                case_key,
                normalized_type.value,
                force,
                audit_context=get_log_context(),
            )
            queued = self._queue_status(
                username=username,
                case_key=case_key,
                artifact_type=normalized_type,
                job_id=result.id,
                reused=False,
            )
        else:
            job_id = f"inprocess-rda-{normalized_type.value}-{case_key}-{int(datetime.now().timestamp())}"
            queued = self._queue_status(
                username=username,
                case_key=case_key,
                artifact_type=normalized_type,
                job_id=job_id,
                reused=False,
            )
            self._spawn_inprocess_job(
                username=username,
                case_key=case_key,
                artifact_type=normalized_type,
                force=force,
                job_id=job_id,
            )

        self.audit_logger.business_event(
            event_type="rda.queued",
            action="queue_case_rda",
            outcome="success",
            service="rda_service",
            resource={
                "case_key": case_key,
                "artifact_type": normalized_type.value,
                "job_id": queued.job_id,
                "status": self._job_state_value(queued.status),
                "reused": queued.reused,
            },
        )
        return queued

    def generate_case_rda(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
        *,
        force: bool = False,
    ):
        normalized_type = self._normalize_artifact_type(artifact_type)
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        command = GenerateRdaCommand(
            username=username,
            case_key=case_key,
            artifact_type=normalized_type,
            force=force,
            persist=True,
        )
        result = self.generate_case_rda_use_case.execute(command)
        self.audit_logger.business_event(
            event_type="rda.generated" if result.source == "generated" else "rda.cache_hit",
            action="generate_case_rda",
            outcome="success",
            service="rda_service",
            resource={
                "case_key": case_key,
                "artifact_type": normalized_type.value,
                "version": result.version,
                "source": result.source,
                "completeness_status": result.completeness_status,
                "force": bool(force),
            },
        )
        return result

    def run_generation_job(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
        *,
        force: bool = False,
        job_id: str = "",
        audit_context: dict[str, object] | None = None,
    ) -> None:
        normalized_type = self._normalize_artifact_type(artifact_type)
        set_log_context(dict(audit_context or {}))
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        self.job_status_repository.upsert_status(
            username=username,
            case_key=case_key,
            artifact_type=normalized_type,
            status=RdaJobState.PROCESSING.value,
            job_id=job_id,
            reused=False,
        )
        self.audit_logger.business_event(
            event_type="rda.processing",
            action="run_generation_job",
            outcome="success",
            service="rda_service",
            resource={"case_key": case_key, "artifact_type": normalized_type.value, "job_id": job_id},
        )
        try:
            result = self.generate_case_rda(
                username,
                case_key,
                normalized_type,
                force=force,
            )
            self.job_status_repository.upsert_status(
                username=username,
                case_key=case_key,
                artifact_type=normalized_type,
                status=RdaJobState.COMPLETED.value,
                job_id=job_id,
                source_context_fingerprint=result.source_context_fingerprint,
                resolved_version=result.version,
                reused=result.source == "cache",
            )
        except Exception as exc:
            self.job_status_repository.upsert_status(
                username=username,
                case_key=case_key,
                artifact_type=normalized_type,
                status=RdaJobState.FAILED.value,
                job_id=job_id,
                error=str(exc),
                reused=False,
            )
            self.audit_logger.business_event(
                event_type="rda.failed",
                action="run_generation_job",
                outcome="error",
                service="rda_service",
                resource={"case_key": case_key, "artifact_type": normalized_type.value, "job_id": job_id},
                error={"class": exc.__class__.__name__, "message": str(exc)},
            )
            raise
        finally:
            clear_log_context()

    def get_latest_case_rda(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
    ):
        normalized_type = self._normalize_artifact_type(artifact_type)
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        result = self.generate_case_rda_use_case.get_latest(username, case_key, normalized_type)
        if result is not None:
            self.audit_logger.business_event(
                event_type="rda.read",
                action="get_latest_case_rda",
                outcome="success",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                    "version": result.version,
                    "completeness_status": result.completeness_status,
                },
            )
        else:
            self.audit_logger.business_event(
                event_type="rda.read_missing",
                action="get_latest_case_rda",
                outcome="not_found",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                },
            )
        return result

    def get_case_rda_status(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
    ) -> RdaStatusResult | None:
        normalized_type = self._normalize_artifact_type(artifact_type)
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        result = self.generate_case_rda_use_case.get_status(username, case_key, normalized_type)
        if result is not None:
            self.audit_logger.business_event(
                event_type="rda.status",
                action="get_case_rda_status",
                outcome="success",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                    "status": self._job_state_value(result.status),
                    "job_id": result.job_id,
                    "resolved_version": result.resolved_version,
                    "source": result.source,
                },
            )
        else:
            self.audit_logger.business_event(
                event_type="rda.status_missing",
                action="get_case_rda_status",
                outcome="not_found",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                },
            )
        return result

    def get_case_rda_validation(
        self,
        username: str,
        case_key: str,
        artifact_type: str | RdaArtifactType,
    ) -> RdaValidationReviewResult | None:
        normalized_type = self._normalize_artifact_type(artifact_type)
        bind_log_context(username=username, case_key=case_key, document_type=f"rda:{normalized_type.value}")
        artifact = self.generate_case_rda_use_case.get_latest(username, case_key, normalized_type)
        status = self.generate_case_rda_use_case.get_status(username, case_key, normalized_type)
        validation = (
            self.validate_rda_fhir_projection_use_case.execute(artifact)
            if artifact is not None
            else None
        )

        if artifact is None and status is None:
            self.audit_logger.business_event(
                event_type="rda.validation_missing",
                action="get_case_rda_validation",
                outcome="not_found",
                service="rda_service",
                resource={
                    "case_key": case_key,
                    "artifact_type": normalized_type.value,
                },
            )
            return None

        review = RdaValidationReviewResult(
            case_key=case_key,
            artifact_type=normalized_type,
            artifact=artifact,
            status=status,
            validation=validation,
        )
        self.audit_logger.business_event(
            event_type="rda.validation_read",
            action="get_case_rda_validation",
            outcome="success",
            service="rda_service",
            resource={
                "case_key": case_key,
                "artifact_type": normalized_type.value,
                "has_artifact": artifact is not None,
                "has_status": status is not None,
                "status": self._job_state_value(status.status) if status is not None else "",
                "resolved_version": status.resolved_version if status is not None else None,
                "error_count": validation.error_count if validation is not None else 0,
                "warning_count": validation.warning_count if validation is not None else 0,
            },
        )
        return review
