"""Casos de uso para la carga individual asíncrona."""

from __future__ import annotations

import hashlib
import logging
from dataclasses import dataclass, field
from datetime import tzinfo
from typing import Any, Literal, Protocol
from uuid import uuid4

from app.batch_processing.application.utils import now_iso
from app.batch_processing.domain.ports import (
    CaseAssociationService,
    DocumentClassifier,
    TextExtractor,
)
from app.clinical_pipeline.domain.errors import HistoriaSummaryCoverageError
from app.clinical_pipeline.domain.models import ProcessingFailureMetadata
from app.core.logging import bind_log_context, get_audit_logger
from app.individual_ingestion.application.models import (
    IndividualUploadAcceptedView,
    IndividualUploadCreatePayload,
    IndividualUploadSessionView,
    IndividualUploadStatusView,
    IndividualUploadUpdatePayload,
)
from app.individual_ingestion.application.support_identity_matching import SupportIdentityMatchEngine
from app.individual_ingestion.domain.models import (
    ACTIVE_SESSION_STATUSES,
    PENDING_ACTION_CONFIRM_CASE_IDENTITY,
    PENDING_ACTION_CONFIRM_CASE_OVERRIDE,
    PENDING_ACTION_CONFIRM_TYPE_OVERRIDE,
    SESSION_STATUS_ACTIVE,
    SESSION_STATUS_BLOCKED,
    SESSION_STATUS_CANCELLED,
    SESSION_STATUS_COMPLETED,
    SUPPORTED_INDIVIDUAL_SUPPORT_TYPES,
    TERMINAL_UPLOAD_STATUSES,
    UPLOAD_STATUS_BLOCKED_BASE_FAILED,
    UPLOAD_STATUS_CANCELADO,
    UPLOAD_STATUS_COMPLETADO,
    UPLOAD_STATUS_ELIMINADO,
    UPLOAD_STATUS_ESPERANDO_CONFIRMACION,
    UPLOAD_STATUS_FALLIDO,
    UPLOAD_STATUS_LISTO_PARA_PROCESAR,
    UPLOAD_STATUS_PRECHECK_EN_COLA,
    UPLOAD_STATUS_PRECHECK_PROCESANDO,
    UPLOAD_STATUS_PROCESANDO,
    UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
)
from app.llm import LLMProviderError
from app.services.clinical_document_service import ClinicalDocumentRequest


logger = logging.getLogger(__name__)
audit_logger = get_audit_logger()


class IndividualUploadDispatcher(Protocol):
    def dispatch_precheck(self, upload_id: str) -> None: ...

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


class IndividualClinicalDocumentService(Protocol):
    def resolve_case_identity(self, request: ClinicalDocumentRequest) -> Any: ...

    def get_user_case_context(self, username: str, case_key: str) -> dict[str, Any]: ...

    def process_and_persist(self, request: ClinicalDocumentRequest) -> dict[str, Any]: ...


def _compute_file_hash(contents: bytes) -> str:
    return hashlib.sha256(contents).hexdigest()


def _impact_summary(selected_type: str, detected_type: str) -> list[str]:
    if not selected_type or not detected_type or selected_type == detected_type:
        return []
    impact_map = {
        "historia_clinica": [
            "Puede alterar la base clínica del caso y la extracción de diagnósticos, procedimientos y medicamentos.",
            "Puede afectar la elegibilidad y consistencia de epicrisis, RIPS y RDA.",
        ],
        "factura": [
            "Puede alterar la extracción administrativa y la consistencia de cobros frente al caso clínico.",
            "Puede afectar validaciones de soportes y hallazgos posteriores.",
        ],
        "quirurgico": [
            "Puede alterar procedimientos extraídos y codificación CUPS/CIE-10 quirúrgica.",
            "Puede afectar la consistencia de artefactos derivados del caso.",
        ],
        "laboratorio": [
            "Puede alterar ayudas diagnósticas y hallazgos consolidados del caso.",
        ],
        "radiologia": [
            "Puede alterar ayudas diagnósticas y hallazgos consolidados del caso.",
        ],
        "generico": [
            "Puede ocultar reglas específicas del tipo documental real y afectar trazabilidad.",
        ],
    }
    return list(impact_map.get(selected_type, []))


def _is_terminal_status(status: str) -> bool:
    return status in TERMINAL_UPLOAD_STATUSES


def _is_deleted_status(status: str) -> bool:
    return status == UPLOAD_STATUS_ELIMINADO


def _conditional_status_update(
    repository: Any,
    upload_id: str,
    expected_statuses: set[str],
    payload: dict[str, Any],
) -> bool:
    update_method = getattr(repository, "update_upload_if_status", None)
    if callable(update_method):
        return bool(update_method(upload_id, expected_statuses, payload))
    repository.update_upload(upload_id, payload)
    return True


def _is_base_terminal_failure(status: str) -> bool:
    return status in {UPLOAD_STATUS_CANCELADO, UPLOAD_STATUS_FALLIDO, UPLOAD_STATUS_BLOCKED_BASE_FAILED}


def _is_story_type(document_type: str) -> bool:
    return _normalize_type(document_type) == "historia_clinica"


def _has_retry_artifacts(record: dict[str, Any]) -> bool:
    return bool(
        str(record.get("extracted_text") or "").strip() or str(record.get("stored_path") or "").strip()
    )


def _legacy_retryable_failure(record: dict[str, Any]) -> ProcessingFailureMetadata | None:
    error = str(record.get("error") or "").strip()
    if error != "No fue posible resumir todos los bloques de la historia clínica.":
        return None
    return ProcessingFailureMetadata(
        code="external_provider_temporarily_unavailable",
        source="external_provider",
        stage="historia_chunk_summary",
        retryable=True,
        provider="gemini",
        error_kind="transient",
        public_title="Servicio externo temporalmente no disponible",
        public_message=(
            "El servicio externo de IA está temporalmente no disponible. "
            "El documento fue recibido correctamente y permanece guardado; "
            "no necesitas cargarlo nuevamente."
        ),
    )


def _project_retryability(record: dict[str, Any]) -> dict[str, Any]:
    projected = dict(record)
    failure = dict(projected.get("failure") or {})
    if not failure:
        legacy = _legacy_retryable_failure(projected)
        if legacy:
            failure = legacy.to_document()
    retryable = bool(failure.get("retryable")) and _has_retry_artifacts(projected)
    projected["failure"] = failure
    projected["can_retry"] = bool(
        str(projected.get("status") or "") == UPLOAD_STATUS_FALLIDO
        and not str(projected.get("clinical_document_id") or "").strip()
        and retryable
    )
    return projected


def _failure_from_exception(exc: Exception) -> ProcessingFailureMetadata:
    if isinstance(exc, HistoriaSummaryCoverageError):
        return exc.to_failure_metadata()
    if isinstance(exc, LLMProviderError):
        external = exc.kind.value not in {"invalid_configuration", "unsupported_capability"}
        external_temporary = external and exc.retryable
        return ProcessingFailureMetadata(
            code="external_provider_temporarily_unavailable"
            if external_temporary
            else "clinical_processing_failed",
            source="external_provider" if external else "internal",
            stage="clinical_document_extract",
            retryable=bool(external and exc.retryable),
            provider=exc.provider,
            model=exc.model or "",
            error_kind=exc.kind.value,
            public_title=(
                "Servicio externo temporalmente no disponible"
                if external_temporary
                else "Fallo del servicio externo de IA"
                if external
                else "No fue posible completar el procesamiento"
            ),
            public_message=(
                "El servicio externo de IA está temporalmente no disponible. "
                "El documento fue recibido correctamente y permanece guardado; "
                "no necesitas cargarlo nuevamente."
                if external_temporary
                else "El servicio externo de IA no pudo completar el procesamiento clínico."
                if external
                else "No fue posible completar el procesamiento clínico."
            ),
        )
    return ProcessingFailureMetadata(
        code="clinical_processing_failed",
        source="internal",
        stage="materialization",
        retryable=False,
        error_kind=exc.__class__.__name__,
        public_title="No fue posible completar el procesamiento",
        public_message="Ocurrió un error interno durante el procesamiento clínico.",
    )


def _precheck_failure_from_exception(exc: Exception) -> ProcessingFailureMetadata:
    message = str(exc).strip()
    if "no contiene texto extraíble" in message.lower():
        return ProcessingFailureMetadata(
            code="document_without_extractable_text",
            source="document",
            stage="pdf_extraction",
            retryable=False,
            error_kind=exc.__class__.__name__,
            public_title="El documento no puede procesarse automáticamente",
            public_message=(
                "El PDF no contiene texto extraíble. Selecciona otro archivo o una versión con texto "
                "reconocible."
            ),
        )
    return ProcessingFailureMetadata(
        code="precheck_failed",
        source="internal",
        stage="precheck",
        retryable=False,
        error_kind=exc.__class__.__name__,
        public_title="No fue posible revisar el documento",
        public_message="Ocurrió un error interno durante la revisión inicial del documento.",
    )


def _create_session_id() -> str:
    return f"manual-{uuid4().hex}"


def _normalize_type(value: Any) -> str:
    return str(value or "").strip()


def _resolve_effective_document_type(
    *,
    selected_type: str,
    detected_type: str,
    requested_type: str = "",
) -> str:
    normalized_selected = _normalize_type(selected_type)
    normalized_detected = _normalize_type(detected_type)
    normalized_requested = _normalize_type(requested_type)
    if normalized_requested and normalized_requested in {normalized_selected, normalized_detected}:
        return normalized_requested
    if normalized_selected == "historia_clinica":
        return normalized_selected
    return normalized_detected or normalized_selected or "generico"


def _sanitize_provisional_identity_value(value: Any, *, unknown_tokens: set[str] | None = None) -> str:
    normalized = _normalize_type(value)
    if not normalized:
        return ""
    lowered = normalized.lower()
    if unknown_tokens and lowered in unknown_tokens:
        return ""
    return normalized


def _normalize_redacted_identity_fields(
    value: Any,
) -> list[Literal["patient_id", "patient_name"]]:
    normalized: list[Literal["patient_id", "patient_name"]] = []
    for item in value or []:
        if (item == "patient_id" or item == "patient_name") and item not in normalized:
            normalized.append(item)
    return normalized


def _serialize_session(session: dict[str, Any] | None) -> dict[str, Any] | None:
    if not session:
        return None
    uploads = []
    for item in session.get("uploads") or []:
        upload_id = str(item.get("_id") or "").strip()
        if not upload_id:
            continue
        uploads.append(
            IndividualUploadStatusView.from_record(
                upload_id,
                _project_retryability(item),
            ).to_document()
        )
    return IndividualUploadSessionView(
        session_id=str(session.get("session_id") or "").strip(),
        active_case_key=str(session.get("active_case_key") or "").strip(),
        base_upload_id=str(session.get("base_upload_id") or "").strip(),
        session_status=str(session.get("session_status") or "").strip(),
        base_ready=bool(session.get("base_ready")),
        blocked_by_base_failure=bool(session.get("blocked_by_base_failure")),
        updated_at=str(session.get("updated_at") or "").strip(),
        uploads=uploads,
    ).to_document()


@dataclass
class CreateIndividualUploadUseCase:
    repository: Any
    file_store: Any
    dispatcher: IndividualUploadDispatcher
    colombia_tz: tzinfo

    def execute(
        self,
        *,
        filename: str,
        contents: bytes,
        username: str,
        selected_document_type: str,
        provided_case_key: str = "",
        session_id: str = "",
        lock_case_selection: bool = False,
    ) -> dict[str, Any]:
        normalized_filename = str(filename or "").strip()
        normalized_type = str(selected_document_type or "").strip()
        normalized_case_key = str(provided_case_key or "").strip()
        normalized_session_id = str(session_id or "").strip()
        if not normalized_filename or not normalized_filename.lower().endswith(".pdf"):
            raise ValueError("Solo se permiten archivos PDF.")
        if not contents:
            raise ValueError("El archivo PDF está vacío.")
        if normalized_type == "historia_clinica" and normalized_case_key:
            raise ValueError("La historia clínica individual no debe llegar con case_key preexistente.")
        if normalized_type != "historia_clinica":
            if normalized_type not in SUPPORTED_INDIVIDUAL_SUPPORT_TYPES:
                raise ValueError("Tipo de soporte individual no soportado.")
            if not normalized_case_key:
                raise ValueError("Debes seleccionar un caso existente para adjuntar un soporte.")

        file_hash = _compute_file_hash(contents)
        duplicate = self.repository.find_latest_by_hash(username, file_hash)
        if duplicate and _is_deleted_status(str(duplicate.get("status") or "")):
            duplicate = None
        if duplicate and not _is_terminal_status(str(duplicate.get("status") or "")):
            upload_id = str(duplicate.get("_id") or "")
            return IndividualUploadAcceptedView(
                upload_id=upload_id,
                status=str(duplicate.get("status") or UPLOAD_STATUS_PRECHECK_EN_COLA),
                created_at=str(duplicate.get("created_at") or ""),
                poll_url=f"/api/cargue-individual/{upload_id}",
                duplicate=True,
                message="El documento ya está en curso.",
            ).to_document()
        if duplicate and str(duplicate.get("clinical_document_id") or "").strip():
            upload_id = str(duplicate.get("_id") or "")
            return IndividualUploadAcceptedView(
                upload_id=upload_id,
                status=str(duplicate.get("status") or UPLOAD_STATUS_COMPLETADO),
                created_at=str(duplicate.get("created_at") or ""),
                poll_url=f"/api/cargue-individual/{upload_id}",
                duplicate=True,
                message="El documento ya fue cargado previamente.",
            ).to_document()
        if duplicate:
            retryable_duplicate = _project_retryability(duplicate)
            if retryable_duplicate.get("can_retry"):
                upload_id = str(duplicate.get("_id") or "")
                return IndividualUploadAcceptedView(
                    upload_id=upload_id,
                    status=UPLOAD_STATUS_FALLIDO,
                    created_at=str(duplicate.get("created_at") or ""),
                    poll_url=f"/api/cargue-individual/{upload_id}",
                    duplicate=True,
                    message=(
                        "El documento ya está registrado y conserva una falla recuperable. "
                        "Reintenta el procesamiento sin cargar otro archivo."
                    ),
                ).to_document()

        created_at = now_iso(self.colombia_tz)
        requested_session = (
            self.repository.get_session(username, normalized_session_id) if normalized_session_id else None
        ) or {}
        requested_session_status = str(requested_session.get("session_status") or "").strip()
        use_existing_session = bool(
            not _is_story_type(normalized_type)
            and normalized_session_id
            and requested_session
            and requested_session_status in ACTIVE_SESSION_STATUSES
        )
        effective_session_id = normalized_session_id if use_existing_session else ""
        base_upload_id = ""
        depends_on_upload_id = ""
        session_case_key = ""
        session_status = ""
        if _is_story_type(normalized_type):
            effective_session_id = _create_session_id()
        elif use_existing_session:
            if bool(requested_session.get("blocked_by_base_failure")):
                raise ValueError(
                    "La sesión manual seleccionada está bloqueada por falla de la historia base."
                )
            base_upload_id = str(requested_session.get("base_upload_id") or "").strip()
            depends_on_upload_id = base_upload_id
            session_case_key = str(requested_session.get("active_case_key") or normalized_case_key).strip()
            session_status = requested_session_status or SESSION_STATUS_ACTIVE

        payload = IndividualUploadCreatePayload(
            usuario=username,
            original_name=normalized_filename,
            selected_document_type=normalized_type,
            provided_case_key=normalized_case_key,
            source_file_hash=file_hash,
            stored_path="",
            file_size=len(contents),
            created_at=created_at,
            updated_at=created_at,
            effective_document_type=normalized_type,
            session_id=effective_session_id,
            session_status=SESSION_STATUS_ACTIVE if effective_session_id else session_status,
            session_origin="manual_upload_dialog" if effective_session_id else "",
            base_upload_id=base_upload_id,
            depends_on_upload_id=depends_on_upload_id,
            session_case_key=session_case_key,
            case_selection_locked=bool(lock_case_selection and normalized_case_key),
        )
        upload_id = self.repository.create_upload(payload.to_document())
        if _is_story_type(normalized_type):
            self.repository.update_upload(
                upload_id,
                IndividualUploadUpdatePayload(
                    base_upload_id=upload_id,
                    updated_at=created_at,
                ).to_document(),
            )
        stored_path = self.file_store.save_file(upload_id, normalized_filename, contents)
        self.repository.update_upload(
            upload_id,
            IndividualUploadUpdatePayload(stored_path=stored_path, updated_at=created_at).to_document(),
        )
        self.dispatcher.dispatch_precheck(upload_id)
        return IndividualUploadAcceptedView(
            upload_id=upload_id,
            status=UPLOAD_STATUS_PRECHECK_EN_COLA,
            created_at=created_at,
            poll_url=f"/api/cargue-individual/{upload_id}",
            duplicate=False,
            message="Documento aceptado para precheck.",
        ).to_document()


@dataclass
class GetIndividualUploadStatusUseCase:
    repository: Any

    def execute(self, upload_id: str, *, username: str) -> dict[str, Any] | None:
        record = self.repository.get_upload(upload_id)
        if not record or str(record.get("usuario") or "").strip() != username:
            return None
        return IndividualUploadStatusView.from_record(
            upload_id,
            _project_retryability(record),
        ).to_document()


@dataclass
class GetLatestPendingIndividualUploadUseCase:
    repository: Any

    def execute(self, username: str) -> dict[str, Any] | None:
        record = self.repository.find_latest_active(username)
        if not record:
            return None
        upload_id = str(record.get("_id") or "")
        return IndividualUploadStatusView.from_record(upload_id, record).to_document()


@dataclass
class GetActiveIndividualUploadSessionUseCase:
    repository: Any

    def execute(self, username: str) -> dict[str, Any] | None:
        session = self.repository.get_active_session(username)
        return _serialize_session(session)


@dataclass
class GetIndividualUploadSessionUseCase:
    repository: Any

    def execute(self, username: str, *, session_id: str) -> dict[str, Any] | None:
        session = self.repository.get_session(username, session_id)
        return _serialize_session(session)


@dataclass
class ListIndividualUploadSessionsUseCase:
    repository: Any

    def execute(self, username: str, *, include_terminal: bool = False) -> list[dict[str, Any]]:
        sessions = self.repository.list_sessions(username, include_terminal=include_terminal)
        return [serialized for item in sessions if (serialized := _serialize_session(item)) is not None]


@dataclass
class CloseIndividualUploadSessionUseCase:
    repository: Any
    colombia_tz: tzinfo

    def execute(self, *, username: str, session_id: str) -> dict[str, Any]:
        session = self.repository.get_session(username, session_id) or {}
        if not session:
            raise ValueError("Sesión manual no encontrada.")
        uploads = list(session.get("uploads") or [])
        if not uploads:
            raise ValueError("La sesión manual no contiene cargas.")
        non_terminal_uploads = [
            item for item in uploads if str(item.get("status") or "").strip() not in TERMINAL_UPLOAD_STATUSES
        ]
        if non_terminal_uploads:
            raise ValueError("Debes cancelar o completar todas las cargas antes de cerrar la sesión manual.")
        updated_at = now_iso(self.colombia_tz)
        for item in uploads:
            upload_id = str(item.get("_id") or "").strip()
            if not upload_id:
                continue
            self.repository.update_upload(
                upload_id,
                IndividualUploadUpdatePayload(
                    session_status=SESSION_STATUS_COMPLETED,
                    updated_at=updated_at,
                ).to_document(),
            )
        return self.repository.get_session(username, session_id) or {
            "session_id": session_id,
            "session_status": SESSION_STATUS_COMPLETED,
            "uploads": [],
        }


@dataclass
class RunIndividualPrecheckUseCase:
    repository: Any
    file_store: Any
    text_extractor: TextExtractor
    classifier: DocumentClassifier
    clinical_document_service: IndividualClinicalDocumentService
    case_association_service: CaseAssociationService
    dispatcher: IndividualUploadDispatcher
    colombia_tz: tzinfo
    support_identity_match_engine: SupportIdentityMatchEngine = field(
        default_factory=SupportIdentityMatchEngine
    )

    def execute(self, upload_id: str) -> dict[str, Any]:
        record = self.repository.get_upload(upload_id) or {}
        if not record:
            return {}
        bind_log_context(
            username=record.get("usuario", ""),
            file_id=upload_id,
            document_type=record.get("selected_document_type", ""),
        )
        current_status = str(record.get("status") or "")
        if current_status not in {UPLOAD_STATUS_PRECHECK_EN_COLA, UPLOAD_STATUS_PRECHECK_PROCESANDO}:
            return record
        claimed = _conditional_status_update(
            self.repository,
            upload_id,
            {UPLOAD_STATUS_PRECHECK_EN_COLA, UPLOAD_STATUS_PRECHECK_PROCESANDO},
            IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_PRECHECK_PROCESANDO,
                error="",
                updated_at=now_iso(self.colombia_tz),
            ).to_document(),
        )
        if not claimed:
            return self.repository.get_upload(upload_id) or {}
        record = self.repository.get_upload(upload_id) or {}
        if _is_deleted_status(str(record.get("status") or "")):
            return record

        try:
            stored_path = str(record.get("stored_path") or "")
            extract_method = getattr(self.text_extractor, "extract", None)
            if callable(extract_method):
                extraction_result = extract_method(stored_path)
                raw_text = extraction_result.text
                extraction_metadata = extraction_result.metadata.to_document()
            else:
                raw_text = self.text_extractor.extract_text(stored_path)
                extraction_metadata = {
                    "character_count": len(raw_text),
                    "extractor": self.text_extractor.__class__.__name__,
                    "warnings": [],
                }
            if not raw_text.strip():
                raise ValueError("El PDF no contiene texto extraíble y requiere revisión manual.")
            selected_type = str(record.get("selected_document_type") or "").strip()
            original_name = str(record.get("original_name") or "")
            inspect_method = getattr(self.classifier, "inspect", None)
            if callable(inspect_method):
                classification = inspect_method(original_name, raw_text)
                detected_type = classification.document_type
                classification_decision = classification.to_document()
            else:
                detected_type = self.classifier.classify(original_name, raw_text)
                classification_decision = {
                    "document_type": detected_type,
                    "title": detected_type.replace("_", " "),
                    "confidence": 0.0,
                    "reasons": ["legacy_classifier"],
                    "scores": {},
                }
            pending_actions: list[str] = []
            warnings: list[str] = []
            contradictions: list[str] = []
            provisional_case_number = ""
            provisional_patient_id = ""
            provisional_patient_name = ""
            redacted_identity_fields: list[Literal["patient_id", "patient_name"]] = []
            case_key = str(record.get("provided_case_key") or "").strip()

            if selected_type == "historia_clinica":
                resolution = self.clinical_document_service.resolve_case_identity(
                    ClinicalDocumentRequest(
                        raw_text=raw_text,
                        detected_type="historia_clinica",
                        username=str(record.get("usuario") or ""),
                        original_name=str(record.get("original_name") or ""),
                    )
                )
                provisional_case_number = resolution.case_number
                provisional_patient_id = _sanitize_provisional_identity_value(resolution.patient_id)
                provisional_patient_name = _sanitize_provisional_identity_value(
                    resolution.patient_name,
                    unknown_tokens={"desconocido"},
                )
                redacted_identity_fields = _normalize_redacted_identity_fields(
                    getattr(resolution, "redacted_identity_fields", [])
                )
                case_key = str(getattr(resolution, "case_key", "") or "").strip()
                warnings.extend(list(resolution.review_messages or []))
                if not provisional_patient_id:
                    message = (
                        "La identificación del paciente está censurada por políticas de protección de datos."
                        if "patient_id" in redacted_identity_fields
                        else "No fue posible identificar con confianza la identificación del paciente."
                    )
                    if message not in warnings:
                        warnings.append(message)
                if not provisional_patient_name:
                    message = (
                        "El nombre del paciente está censurado por políticas de protección de datos."
                        if "patient_name" in redacted_identity_fields
                        else "No fue posible identificar con confianza el nombre del paciente."
                    )
                    if message not in warnings:
                        warnings.append(message)
                pending_actions.append(PENDING_ACTION_CONFIRM_CASE_IDENTITY)
            else:
                case_context = self.clinical_document_service.get_user_case_context(
                    str(record.get("usuario") or ""),
                    case_key,
                )
                if not case_context:
                    base_upload_id = str(
                        record.get("depends_on_upload_id") or record.get("base_upload_id") or ""
                    ).strip()
                    base_upload = self.repository.get_upload(base_upload_id) or {}
                    if base_upload:
                        case_context = {
                            "case_key": str(
                                base_upload.get("case_key")
                                or base_upload.get("session_case_key")
                                or record.get("provided_case_key")
                                or case_key
                            ).strip(),
                            "case_number": str(
                                base_upload.get("confirmed_case_number")
                                or base_upload.get("provisional_case_number")
                                or ""
                            ).strip(),
                            "patient_id": str(
                                base_upload.get("confirmed_patient_id")
                                or base_upload.get("provisional_patient_id")
                                or ""
                            ).strip(),
                            "patient_name": str(
                                base_upload.get("confirmed_patient_name")
                                or base_upload.get("provisional_patient_name")
                                or ""
                            ).strip(),
                        }
                if not case_context:
                    raise ValueError("No se encontró el caso destino para el soporte.")
                signals = self.case_association_service.extract_signals(
                    str(record.get("original_name") or ""),
                    raw_text,
                    detected_type,
                )
                identity_decision = self.support_identity_match_engine.evaluate(
                    filename=str(record.get("original_name") or ""),
                    raw_text=raw_text,
                    detected_type=detected_type,
                    preferred_case_context=case_context,
                    extracted_signals=signals,
                )
                contradictions = list(identity_decision.contradictions)
                warnings.extend(identity_decision.warnings)
                if contradictions:
                    pending_actions.append(PENDING_ACTION_CONFIRM_CASE_OVERRIDE)

            if detected_type != selected_type:
                pending_actions.append(PENDING_ACTION_CONFIRM_TYPE_OVERRIDE)
                warnings.append(
                    f"El documento parece ser '{detected_type}' pero fue cargado como '{selected_type}'."
                )

            next_status = (
                UPLOAD_STATUS_ESPERANDO_CONFIRMACION if pending_actions else UPLOAD_STATUS_LISTO_PARA_PROCESAR
            )
            impact_summary = _impact_summary(selected_type, detected_type)
            effective_document_type = _resolve_effective_document_type(
                selected_type=selected_type,
                detected_type=detected_type,
            )
            updated_payload = IndividualUploadUpdatePayload(
                status=next_status,
                detected_document_type=detected_type,
                effective_document_type=effective_document_type,
                extracted_text=raw_text,
                extraction_metadata=extraction_metadata,
                classification_decision=classification_decision,
                pending_actions=pending_actions,
                warnings=warnings,
                contradictions=contradictions,
                impact_summary=impact_summary,
                provisional_case_number=provisional_case_number,
                provisional_patient_id=provisional_patient_id,
                provisional_patient_name=provisional_patient_name,
                redacted_identity_fields=redacted_identity_fields,
                case_key=case_key,
                session_case_key=case_key,
                precheck_completed_at=now_iso(self.colombia_tz),
                updated_at=now_iso(self.colombia_tz),
            ).to_document()
            updated = _conditional_status_update(
                self.repository,
                upload_id,
                {UPLOAD_STATUS_PRECHECK_PROCESANDO},
                updated_payload,
            )
            if not updated:
                return self.repository.get_upload(upload_id) or {}
            updated_record = self.repository.get_upload(upload_id) or {}
            if next_status == UPLOAD_STATUS_LISTO_PARA_PROCESAR and not _is_deleted_status(
                str(updated_record.get("status") or "")
            ):
                self.dispatcher.dispatch_materialization(upload_id)
            return updated_record
        except Exception as exc:
            logger.exception("Error en precheck de carga individual %s", upload_id)
            failure = _precheck_failure_from_exception(exc)
            _conditional_status_update(
                self.repository,
                upload_id,
                {UPLOAD_STATUS_PRECHECK_PROCESANDO},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_FALLIDO,
                    failure=failure.to_document(),
                    can_retry=False,
                    error=failure.public_message,
                    failed_at=now_iso(self.colombia_tz),
                    updated_at=now_iso(self.colombia_tz),
                ).to_document(),
            )
            return self.repository.get_upload(upload_id) or {}


@dataclass
class ConfirmIndividualUploadUseCase:
    repository: Any
    dispatcher: IndividualUploadDispatcher
    colombia_tz: tzinfo

    def execute(
        self,
        *,
        upload_id: str,
        username: str,
        confirmed_case_number: str = "",
        confirmed_patient_id: str = "",
        confirmed_patient_name: str = "",
        confirmed_effective_document_type: str = "",
    ) -> dict[str, Any]:
        record = self.repository.get_upload(upload_id) or {}
        if not record or str(record.get("usuario") or "").strip() != username:
            raise ValueError("Carga individual no encontrada.")
        if str(record.get("status") or "") != UPLOAD_STATUS_ESPERANDO_CONFIRMACION:
            raise ValueError("La carga no está esperando confirmación.")

        pending_actions = list(record.get("pending_actions") or [])
        normalized_case_number = str(
            confirmed_case_number or record.get("provisional_case_number") or ""
        ).strip()
        normalized_patient_id = str(
            confirmed_patient_id or record.get("provisional_patient_id") or ""
        ).strip()
        normalized_patient_name = str(
            confirmed_patient_name or record.get("provisional_patient_name") or ""
        ).strip()
        redacted_identity_fields = set(record.get("redacted_identity_fields") or [])

        if PENDING_ACTION_CONFIRM_CASE_IDENTITY in pending_actions:
            missing_required_identity = (
                not normalized_case_number
                or (not normalized_patient_id and "patient_id" not in redacted_identity_fields)
                or (not normalized_patient_name and "patient_name" not in redacted_identity_fields)
            )
            if missing_required_identity:
                raise ValueError(
                    "Debes confirmar el número de caso y los datos de identidad no censurados antes de continuar."
                )

        selected_type = str(record.get("selected_document_type") or "").strip()
        detected_type = str(record.get("detected_document_type") or "").strip()
        effective_type = _resolve_effective_document_type(
            selected_type=selected_type,
            detected_type=detected_type,
            requested_type=confirmed_effective_document_type
            or str(record.get("effective_document_type") or ""),
        )
        selected_override_confirmed = bool(
            selected_type
            and detected_type
            and selected_type != detected_type
            and effective_type == selected_type
        )
        case_override_confirmed = PENDING_ACTION_CONFIRM_CASE_OVERRIDE in pending_actions
        override_audit = {}
        if selected_type != detected_type or case_override_confirmed:
            override_audit = {
                "confirmed_by": username,
                "confirmed_at": now_iso(self.colombia_tz),
                "selected_document_type": selected_type,
                "detected_document_type": detected_type,
                "effective_document_type": effective_type,
                "pending_actions": pending_actions,
                "contradictions": list(record.get("contradictions") or []),
                "warnings": list(record.get("warnings") or []),
            }

        self.repository.update_upload(
            upload_id,
            IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_LISTO_PARA_PROCESAR,
                pending_actions=[],
                confirmed_case_number=normalized_case_number,
                confirmed_patient_id=normalized_patient_id,
                confirmed_patient_name=normalized_patient_name,
                case_key=str(record.get("case_key") or ""),
                session_case_key=str(record.get("case_key") or record.get("session_case_key") or ""),
                session_status=SESSION_STATUS_ACTIVE if str(record.get("session_id") or "").strip() else None,
                effective_document_type=effective_type,
                selected_type_override_confirmed=selected_override_confirmed,
                case_override_confirmed=case_override_confirmed,
                override_audit=override_audit,
                updated_at=now_iso(self.colombia_tz),
            ).to_document(),
        )
        self.dispatcher.dispatch_materialization(upload_id)
        return self.repository.get_upload(upload_id) or {}


@dataclass
class RetryIndividualUploadUseCase:
    repository: Any
    dispatcher: IndividualUploadDispatcher
    colombia_tz: tzinfo

    def execute(self, *, upload_id: str, username: str) -> dict[str, Any]:
        record = self.repository.get_upload(upload_id) or {}
        if not record or str(record.get("usuario") or "").strip() != username:
            raise ValueError("Carga individual no encontrada.")

        current_status = str(record.get("status") or "").strip()
        if current_status in {
            UPLOAD_STATUS_PRECHECK_EN_COLA,
            UPLOAD_STATUS_PRECHECK_PROCESANDO,
            UPLOAD_STATUS_LISTO_PARA_PROCESAR,
            UPLOAD_STATUS_PROCESANDO,
            UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
        }:
            return _project_retryability(record)
        if current_status != UPLOAD_STATUS_FALLIDO:
            raise ValueError("La carga no está disponible para reintento.")

        projected = _project_retryability(record)
        if not projected.get("can_retry"):
            raise ValueError("La falla no es recuperable con un reintento.")

        has_extracted_text = bool(str(record.get("extracted_text") or "").strip())
        has_stored_file = bool(str(record.get("stored_path") or "").strip())
        if not has_extracted_text and not has_stored_file:
            raise ValueError("La carga ya no conserva artefactos para reintentar.")

        next_status = (
            UPLOAD_STATUS_LISTO_PARA_PROCESAR if has_extracted_text else UPLOAD_STATUS_PRECHECK_EN_COLA
        )
        retried_at = now_iso(self.colombia_tz)
        payload = IndividualUploadUpdatePayload(
            status=next_status,
            session_status=SESSION_STATUS_ACTIVE
            if str(record.get("session_id") or "").strip()
            else str(record.get("session_status") or ""),
            failure={},
            can_retry=False,
            last_retry_at=retried_at,
            failed_at="",
            error="",
            blocked_by_base_failure=False,
            updated_at=retried_at,
        ).to_document()

        claim_method = getattr(self.repository, "claim_retry", None)
        if callable(claim_method):
            claimed = bool(claim_method(upload_id, username, payload))
        else:
            self.repository.update_upload(
                upload_id,
                {
                    **payload,
                    "retry_count": int(record.get("retry_count") or 0) + 1,
                },
            )
            claimed = True
        if not claimed:
            current = self.repository.get_upload(upload_id) or {}
            current_status = str(current.get("status") or "").strip()
            if current_status in {
                UPLOAD_STATUS_PRECHECK_EN_COLA,
                UPLOAD_STATUS_PRECHECK_PROCESANDO,
                UPLOAD_STATUS_LISTO_PARA_PROCESAR,
                UPLOAD_STATUS_PROCESANDO,
            }:
                return _project_retryability(current)
            raise ValueError("La carga cambió de estado y no pudo reintentarse.")

        if str(record.get("base_upload_id") or "").strip() == upload_id:
            self._restore_dependent_uploads(
                username=username, base_upload_id=upload_id, updated_at=retried_at
            )

        audit_logger.business_event(
            event_type="individual.upload_retried",
            action="retry_upload",
            outcome="success",
            service="individual_ingestion",
            resource={
                "upload_id": upload_id,
                "document_type": record.get("selected_document_type", ""),
            },
            metrics={"retry_count": int(record.get("retry_count") or 0) + 1},
        )
        if next_status == UPLOAD_STATUS_LISTO_PARA_PROCESAR:
            self.dispatcher.dispatch_materialization(upload_id)
        else:
            self.dispatcher.dispatch_precheck(upload_id)
        return _project_retryability(self.repository.get_upload(upload_id) or {})

    def _restore_dependent_uploads(self, *, username: str, base_upload_id: str, updated_at: str) -> None:
        for dependent in self.repository.list_dependent_uploads(username, base_upload_id):
            if str(dependent.get("status") or "").strip() != UPLOAD_STATUS_BLOCKED_BASE_FAILED:
                continue
            dependent_id = str(dependent.get("_id") or "").strip()
            if not dependent_id:
                continue
            _conditional_status_update(
                self.repository,
                dependent_id,
                {str(dependent.get("status") or "").strip()},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
                    session_status=SESSION_STATUS_ACTIVE,
                    waiting_for_base_upload_id=base_upload_id,
                    blocked_by_base_failure=False,
                    failed_at="",
                    error="",
                    updated_at=updated_at,
                ).to_document(),
            )


@dataclass
class CancelIndividualUploadUseCase:
    repository: Any
    file_store: Any
    colombia_tz: tzinfo

    def execute(self, *, upload_id: str, username: str, reason: str = "") -> dict[str, Any]:
        record = self.repository.get_upload(upload_id) or {}
        if not record or str(record.get("usuario") or "").strip() != username:
            raise ValueError("Carga individual no encontrada.")
        current_status = str(record.get("status") or "")
        if current_status in {UPLOAD_STATUS_COMPLETADO, UPLOAD_STATUS_CANCELADO}:
            raise ValueError("La carga ya no se puede cancelar.")
        self.file_store.delete_file(str(record.get("stored_path") or ""))
        audit_logger.business_event(
            event_type="individual.upload_cancelled",
            action="cancel_upload",
            outcome="success",
            service="individual_ingestion",
            resource={"upload_id": upload_id, "document_type": record.get("selected_document_type", "")},
            metrics={"file_size": int(record.get("file_size") or 0)},
            error={"reason": str(reason or "cancelled_by_user")},
        )
        self.repository.update_upload(
            upload_id,
            IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_CANCELADO,
                session_status=SESSION_STATUS_CANCELLED
                if str(record.get("base_upload_id") or "") == upload_id
                else None,
                cancelled_at=now_iso(self.colombia_tz),
                stored_path="",
                extracted_text="",
                pending_actions=[],
                warnings=[],
                contradictions=[],
                impact_summary=[],
                failure={},
                can_retry=False,
                error="",
                audit_only=True,
                updated_at=now_iso(self.colombia_tz),
            ).to_document(),
        )
        if str(record.get("base_upload_id") or "").strip() == upload_id:
            self._block_dependent_uploads(username=username, base_upload_id=upload_id)
        return self.repository.get_upload(upload_id) or {}

    def _block_dependent_uploads(self, *, username: str, base_upload_id: str) -> None:
        for dependent in self.repository.list_dependent_uploads(username, base_upload_id):
            dependent_id = str(dependent.get("_id") or "").strip()
            if not dependent_id:
                continue
            if str(dependent.get("status") or "").strip() in TERMINAL_UPLOAD_STATUSES:
                continue
            self.file_store.delete_file(str(dependent.get("stored_path") or ""))
            _conditional_status_update(
                self.repository,
                dependent_id,
                {str(dependent.get("status") or "").strip()},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_BLOCKED_BASE_FAILED,
                    waiting_for_base_upload_id=base_upload_id,
                    blocked_by_base_failure=True,
                    session_status=SESSION_STATUS_BLOCKED,
                    stored_path="",
                    extracted_text="",
                    error="La historia clínica base de esta sesión fue cancelada o falló.",
                    failed_at=now_iso(self.colombia_tz),
                    updated_at=now_iso(self.colombia_tz),
                ).to_document(),
            )


@dataclass
class RunIndividualMaterializationUseCase:
    repository: Any
    file_store: Any
    clinical_document_service: IndividualClinicalDocumentService
    colombia_tz: tzinfo
    document_cleanup_service: Any | None = None

    def execute(self, upload_id: str) -> dict[str, Any]:
        record = self.repository.get_upload(upload_id) or {}
        if not record:
            return {}
        bind_log_context(
            username=record.get("usuario", ""),
            file_id=upload_id,
            case_key=record.get("provided_case_key", ""),
            document_type=record.get("selected_document_type", ""),
        )
        current_status = str(record.get("status") or "")
        if current_status not in {
            UPLOAD_STATUS_LISTO_PARA_PROCESAR,
            UPLOAD_STATUS_PROCESANDO,
            UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
        }:
            return record
        if _is_deleted_status(current_status):
            return record
        if not _is_story_type(str(record.get("selected_document_type") or "")):
            dependency_result = self._guard_support_dependency(record)
            if dependency_result is not None:
                return dependency_result
        claimed = _conditional_status_update(
            self.repository,
            upload_id,
            {
                UPLOAD_STATUS_LISTO_PARA_PROCESAR,
                UPLOAD_STATUS_PROCESANDO,
                UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
            },
            IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_PROCESANDO,
                materialization_started_at=now_iso(self.colombia_tz),
                failure={},
                can_retry=False,
                updated_at=now_iso(self.colombia_tz),
                error="",
            ).to_document(),
        )
        if not claimed:
            return self.repository.get_upload(upload_id) or {}
        record = self.repository.get_upload(upload_id) or {}
        if _is_deleted_status(str(record.get("status") or "")):
            return record
        try:
            final_case_number = str(
                record.get("confirmed_case_number") or record.get("provisional_case_number") or ""
            ).strip()
            final_patient_id = str(
                record.get("confirmed_patient_id") or record.get("provisional_patient_id") or ""
            ).strip()
            final_patient_name = str(
                record.get("confirmed_patient_name") or record.get("provisional_patient_name") or ""
            ).strip()
            selected_type = str(record.get("selected_document_type") or "").strip()
            detected_type = str(record.get("detected_document_type") or "").strip()
            effective_type = _resolve_effective_document_type(
                selected_type=selected_type,
                detected_type=detected_type,
                requested_type=str(record.get("effective_document_type") or ""),
            )
            result = self.clinical_document_service.process_and_persist(
                ClinicalDocumentRequest(
                    raw_text=str(record.get("extracted_text") or ""),
                    detected_type=effective_type,
                    username=str(record.get("usuario") or ""),
                    original_name=str(record.get("original_name") or ""),
                    case_key=str(record.get("provided_case_key") or ""),
                    case_number=final_case_number,
                    patient_id=final_patient_id,
                    batch_id="",
                    batch_file_id=f"individual:{upload_id}",
                    ingestion_source="individual",
                    provided_patient_name=final_patient_name,
                    provided_patient_id=final_patient_id,
                    provided_case_number=final_case_number,
                    redacted_identity_fields=_normalize_redacted_identity_fields(
                        record.get("redacted_identity_fields") or []
                    ),
                    selected_document_type=selected_type,
                    detected_document_type=detected_type,
                    effective_document_type=effective_type,
                    source_file_hash=str(record.get("source_file_hash") or ""),
                    extraction_metadata=dict(record.get("extraction_metadata") or {}),
                    classification_decision=dict(record.get("classification_decision") or {}),
                    category_override_confirmed=bool(record.get("selected_type_override_confirmed", False)),
                    override_audit=dict(record.get("override_audit") or {}),
                )
            )
            self.file_store.delete_file(str(record.get("stored_path") or ""))
            current = self.repository.get_upload(upload_id) or {}
            if _is_deleted_status(str(current.get("status") or "")):
                self._purge_created_document(record=record, result=result)
                return current
            completed_payload = IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_COMPLETADO,
                session_status=SESSION_STATUS_ACTIVE
                if str(record.get("session_id") or "").strip()
                else str(record.get("session_status") or ""),
                stored_path="",
                extracted_text="",
                clinical_document_id=str(result.get("id_documento") or ""),
                historia_processing=dict(result.get("historia_processing") or {}),
                analysis_quality=dict(result.get("analysis_quality") or {}),
                analysis_provider=str(result.get("analysis_provider") or ""),
                analysis_model_name=str(result.get("analysis_model_name") or ""),
                failure={},
                can_retry=False,
                case_key=str(result.get("case_key") or record.get("provided_case_key") or ""),
                session_case_key=str(result.get("case_key") or record.get("session_case_key") or ""),
                completed_at=now_iso(self.colombia_tz),
                updated_at=now_iso(self.colombia_tz),
            ).to_document()
            completed = _conditional_status_update(
                self.repository,
                upload_id,
                {UPLOAD_STATUS_PROCESANDO},
                completed_payload,
            )
            if not completed:
                current = self.repository.get_upload(upload_id) or {}
                self._purge_created_document(record=record, result=result)
                return current
            updated = self.repository.get_upload(upload_id) or {}
            if str(record.get("base_upload_id") or "").strip() == upload_id:
                self._release_dependent_uploads(updated)
            return updated
        except Exception as exc:
            logger.exception("Error materializando carga individual %s", upload_id)
            failure = _failure_from_exception(exc)
            can_retry = bool(failure.retryable and _has_retry_artifacts(record))
            _conditional_status_update(
                self.repository,
                upload_id,
                {UPLOAD_STATUS_PROCESANDO},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_FALLIDO,
                    session_status=SESSION_STATUS_BLOCKED
                    if str(record.get("base_upload_id") or "").strip() == upload_id
                    else str(record.get("session_status") or ""),
                    error=failure.public_message,
                    failure=failure.to_document(),
                    can_retry=can_retry,
                    failed_at=now_iso(self.colombia_tz),
                    updated_at=now_iso(self.colombia_tz),
            ).to_document(),
            )
            failed = self.repository.get_upload(upload_id) or {}
            if _is_deleted_status(str(failed.get("status") or "")):
                return failed
            if str(record.get("base_upload_id") or "").strip() == upload_id:
                self._block_dependent_uploads(failed)
            return failed

    def _purge_created_document(self, *, record: dict[str, Any], result: dict[str, Any]) -> None:
        document_id = str(result.get("id_documento") or "").strip()
        username = str(record.get("usuario") or "").strip()
        if not document_id or not username or self.document_cleanup_service is None:
            return
        purge = getattr(self.document_cleanup_service, "purge_document_record", None)
        if callable(purge):
            purge(username=username, document_id=document_id)

    def _guard_support_dependency(self, record: dict[str, Any]) -> dict[str, Any] | None:
        username = str(record.get("usuario") or "").strip()
        base_upload_id = str(record.get("depends_on_upload_id") or record.get("base_upload_id") or "").strip()
        if not username or not base_upload_id:
            return None
        base_upload = self.repository.get_upload(base_upload_id) or {}
        base_status = str(base_upload.get("status") or "").strip()
        upload_id = str(record.get("_id") or "").strip()
        if base_status == UPLOAD_STATUS_COMPLETADO:
            if str(record.get("case_key") or "").strip() != str(base_upload.get("case_key") or "").strip():
                updated = _conditional_status_update(
                    self.repository,
                    upload_id,
                    {str(record.get("status") or "").strip()},
                    IndividualUploadUpdatePayload(
                        case_key=str(base_upload.get("case_key") or ""),
                        session_case_key=str(base_upload.get("case_key") or ""),
                        updated_at=now_iso(self.colombia_tz),
                    ).to_document(),
                )
                if not updated:
                    return self.repository.get_upload(upload_id) or {}
            return None
        if _is_base_terminal_failure(base_status):
            updated = _conditional_status_update(
                self.repository,
                upload_id,
                {str(record.get("status") or "").strip()},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_BLOCKED_BASE_FAILED,
                    waiting_for_base_upload_id=base_upload_id,
                    blocked_by_base_failure=True,
                    session_status=SESSION_STATUS_BLOCKED,
                    error="La historia clínica base de esta sesión fue cancelada o falló.",
                    failed_at=now_iso(self.colombia_tz),
                    updated_at=now_iso(self.colombia_tz),
                ).to_document(),
            )
            if not updated:
                return self.repository.get_upload(upload_id) or {}
            return self.repository.get_upload(upload_id) or {}
        updated = _conditional_status_update(
            self.repository,
            upload_id,
            {str(record.get("status") or "").strip()},
            IndividualUploadUpdatePayload(
                status=UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
                waiting_for_base_upload_id=base_upload_id,
                blocked_by_base_failure=False,
                session_status=SESSION_STATUS_ACTIVE,
                updated_at=now_iso(self.colombia_tz),
                error="",
            ).to_document(),
        )
        if not updated:
            return self.repository.get_upload(upload_id) or {}
        return self.repository.get_upload(upload_id) or {}

    def _release_dependent_uploads(self, base_record: dict[str, Any]) -> None:
        username = str(base_record.get("usuario") or "").strip()
        base_upload_id = str(base_record.get("_id") or "").strip()
        if not username or not base_upload_id:
            return
        for dependent in self.repository.list_dependent_uploads(username, base_upload_id):
            dependent_id = str(dependent.get("_id") or "").strip()
            if not dependent_id:
                continue
            if str(dependent.get("status") or "").strip() not in {
                UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
                UPLOAD_STATUS_LISTO_PARA_PROCESAR,
            }:
                continue
            _conditional_status_update(
                self.repository,
                dependent_id,
                {
                    UPLOAD_STATUS_WAITING_BASE_MATERIALIZATION,
                    UPLOAD_STATUS_LISTO_PARA_PROCESAR,
                },
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_LISTO_PARA_PROCESAR,
                    case_key=str(base_record.get("case_key") or ""),
                    session_case_key=str(base_record.get("case_key") or ""),
                    waiting_for_base_upload_id="",
                    blocked_by_base_failure=False,
                    error="",
                    updated_at=now_iso(self.colombia_tz),
                ).to_document(),
            )
            current = self.repository.get_upload(dependent_id) or {}
            if str(current.get("status") or "").strip() == UPLOAD_STATUS_LISTO_PARA_PROCESAR:
                self.execute(dependent_id)

    def _block_dependent_uploads(self, base_record: dict[str, Any]) -> None:
        username = str(base_record.get("usuario") or "").strip()
        base_upload_id = str(base_record.get("_id") or "").strip()
        if not username or not base_upload_id:
            return
        for dependent in self.repository.list_dependent_uploads(username, base_upload_id):
            dependent_id = str(dependent.get("_id") or "").strip()
            if not dependent_id:
                continue
            if str(dependent.get("status") or "").strip() in TERMINAL_UPLOAD_STATUSES:
                continue
            _conditional_status_update(
                self.repository,
                dependent_id,
                {str(dependent.get("status") or "").strip()},
                IndividualUploadUpdatePayload(
                    status=UPLOAD_STATUS_BLOCKED_BASE_FAILED,
                    waiting_for_base_upload_id=base_upload_id,
                    blocked_by_base_failure=True,
                    session_status=SESSION_STATUS_BLOCKED,
                    can_retry=False,
                    error="La historia clínica base falló. Reintenta la historia para continuar.",
                    failed_at=now_iso(self.colombia_tz),
                    updated_at=now_iso(self.colombia_tz),
                ).to_document(),
            )
