"""Runtime del flujo de carga individual aislado del batch masivo."""

from __future__ import annotations

import os
from dataclasses import dataclass
from pathlib import Path
from threading import Thread

from app.batch_processing.infrastructure.heuristic_association import HeuristicCaseAssociationService
from app.batch_processing.infrastructure.heuristic_classifier import HeuristicDocumentClassifier
from app.batch_processing.infrastructure.pdf_text_extractor import PdfTextExtractor
from app.core.logging import get_log_context, set_log_context
from app.individual_ingestion.application.support_identity_matching import SupportIdentityMatchEngine
from app.individual_ingestion.application.use_cases import (
    CancelIndividualUploadUseCase,
    CloseIndividualUploadSessionUseCase,
    ConfirmIndividualUploadUseCase,
    CreateIndividualUploadUseCase,
    GetActiveIndividualUploadSessionUseCase,
    GetIndividualUploadSessionUseCase,
    GetIndividualUploadStatusUseCase,
    GetLatestPendingIndividualUploadUseCase,
    ListIndividualUploadSessionsUseCase,
    RetryIndividualUploadUseCase,
    RunIndividualMaterializationUseCase,
    RunIndividualPrecheckUseCase,
)
from app.individual_ingestion.infrastructure.local_artifacts import LocalIndividualUploadFileStore
from app.individual_ingestion.infrastructure.mongo_repositories import MongoIndividualUploadRepository
from app.services.clinical_document_service import ClinicalDocumentService


class InProcessIndividualDispatcher:
    def __init__(self, *, precheck_runner, materialization_runner) -> None:
        self.precheck_runner = precheck_runner
        self.materialization_runner = materialization_runner

    def dispatch_precheck(self, upload_id: str) -> None:
        self._spawn(self.precheck_runner, upload_id)

    def dispatch_materialization(self, upload_id: str) -> None:
        self._spawn(self.materialization_runner, upload_id)

    def _spawn(self, runner, upload_id: str) -> None:
        audit_context = get_log_context()
        thread = Thread(target=self._safe_run, args=(runner, upload_id, audit_context), daemon=True)
        thread.start()

    def _safe_run(self, runner, upload_id: str, audit_context: dict[str, object]) -> None:
        set_log_context(audit_context)
        runner(upload_id)


class CeleryIndividualDispatcher:
    def dispatch_precheck(self, upload_id: str) -> None:
        from app.batch_processing.celery_app import precheck_individual_upload_job

        precheck_individual_upload_job.apply_async(
            args=[upload_id, get_log_context()], queue="individual_uploads"
        )

    def dispatch_materialization(self, upload_id: str) -> None:
        from app.batch_processing.celery_app import materialize_individual_upload_job

        materialize_individual_upload_job.apply_async(
            args=[upload_id, get_log_context()],
            queue="individual_uploads",
        )


@dataclass
class IndividualIngestionRuntime:
    create_upload: CreateIndividualUploadUseCase
    get_upload_status: GetIndividualUploadStatusUseCase
    get_latest_pending: GetLatestPendingIndividualUploadUseCase
    get_active_session: GetActiveIndividualUploadSessionUseCase
    get_session: GetIndividualUploadSessionUseCase
    list_sessions: ListIndividualUploadSessionsUseCase
    confirm_upload: ConfirmIndividualUploadUseCase
    retry_upload: RetryIndividualUploadUseCase
    cancel_upload: CancelIndividualUploadUseCase
    close_session: CloseIndividualUploadSessionUseCase
    run_precheck: RunIndividualPrecheckUseCase
    run_materialization: RunIndividualMaterializationUseCase
    dispatcher_name: str


def build_individual_ingestion_runtime(
    base_dir: Path,
    colombia_tz,
    *,
    clinical_document_service: ClinicalDocumentService,
    mongo_database=None,
    repository: MongoIndividualUploadRepository | None = None,
    document_cleanup_service=None,
) -> IndividualIngestionRuntime:
    repository = repository or MongoIndividualUploadRepository(database=mongo_database)
    file_store = LocalIndividualUploadFileStore(base_dir)
    attach_file_store = getattr(repository, "attach_file_store", None)
    if callable(attach_file_store):
        attach_file_store(file_store)
    text_extractor = PdfTextExtractor()
    classifier = HeuristicDocumentClassifier()
    association_service = HeuristicCaseAssociationService()
    support_identity_match_engine = SupportIdentityMatchEngine()

    dispatcher_mode = os.getenv("INDIVIDUAL_INGESTION_DISPATCHER", os.getenv("BATCH_DISPATCHER", "inprocess"))
    dispatcher_mode = dispatcher_mode.strip().lower()

    run_materialization = RunIndividualMaterializationUseCase(
        repository=repository,
        file_store=file_store,
        clinical_document_service=clinical_document_service,
        colombia_tz=colombia_tz,
        document_cleanup_service=document_cleanup_service,
    )

    if dispatcher_mode == "celery":
        dispatcher = CeleryIndividualDispatcher()
    else:
        dispatcher = InProcessIndividualDispatcher(
            precheck_runner=None,
            materialization_runner=run_materialization.execute,
        )
        dispatcher_mode = "inprocess"

    run_precheck = RunIndividualPrecheckUseCase(
        repository=repository,
        file_store=file_store,
        text_extractor=text_extractor,
        classifier=classifier,
        clinical_document_service=clinical_document_service,
        case_association_service=association_service,
        dispatcher=dispatcher,
        colombia_tz=colombia_tz,
        support_identity_match_engine=support_identity_match_engine,
    )
    if isinstance(dispatcher, InProcessIndividualDispatcher):
        dispatcher.precheck_runner = run_precheck.execute

    create_upload = CreateIndividualUploadUseCase(
        repository=repository,
        file_store=file_store,
        dispatcher=dispatcher,
        colombia_tz=colombia_tz,
    )
    get_upload_status = GetIndividualUploadStatusUseCase(repository=repository)
    get_latest_pending = GetLatestPendingIndividualUploadUseCase(repository=repository)
    get_active_session = GetActiveIndividualUploadSessionUseCase(repository=repository)
    get_session = GetIndividualUploadSessionUseCase(repository=repository)
    list_sessions = ListIndividualUploadSessionsUseCase(repository=repository)
    confirm_upload = ConfirmIndividualUploadUseCase(
        repository=repository,
        dispatcher=dispatcher,
        colombia_tz=colombia_tz,
    )
    retry_upload = RetryIndividualUploadUseCase(
        repository=repository,
        dispatcher=dispatcher,
        colombia_tz=colombia_tz,
    )
    cancel_upload = CancelIndividualUploadUseCase(
        repository=repository,
        file_store=file_store,
        colombia_tz=colombia_tz,
    )
    close_session = CloseIndividualUploadSessionUseCase(repository=repository, colombia_tz=colombia_tz)

    return IndividualIngestionRuntime(
        create_upload=create_upload,
        get_upload_status=get_upload_status,
        get_latest_pending=get_latest_pending,
        get_active_session=get_active_session,
        get_session=get_session,
        list_sessions=list_sessions,
        confirm_upload=confirm_upload,
        retry_upload=retry_upload,
        cancel_upload=cancel_upload,
        close_session=close_session,
        run_precheck=run_precheck,
        run_materialization=run_materialization,
        dispatcher_name=dispatcher_mode,
    )
