import asyncio
from datetime import datetime, timedelta
from time import perf_counter

from app.core.env import env_float, env_int
from app.core.logger import logger
from app.models.document import IngestRequest
from app.repositories import parse_run_repo
from app.services import ingestion_service


def _retry_backoff_seconds(attempt: int) -> float:
    base = env_float("INGESTION_RETRY_BACKOFF_SECONDS", 1.0)
    multiplier = env_float("INGESTION_RETRY_BACKOFF_MULTIPLIER", 2.0)
    # Deterministic exponential backoff: attempt=1 => base, attempt=2 => base*multiplier, ...
    return max(0.0, base * (multiplier ** max(0, attempt - 1)))


def _is_retryable_processing_error(exc: Exception) -> bool:
    return isinstance(exc, (TimeoutError, ConnectionError))


async def run_ingestion_job(run_id: str, request: IngestRequest) -> None:
    """Execute ingestion with lifecycle tracking, retry policy, and duplicate protection."""
    started_at = perf_counter()
    run = await parse_run_repo.get_run_by_id(run_id)
    if run is None:
        logger.bind(run_id=run_id).warning("Skipping ingestion job for unknown run")
        return

    logger.bind(
        run_id=run_id,
        repo_id=request.repo_id,
        doc_path=request.doc_path,
        initial_status=run.status,
    ).info("Ingestion worker job received")

    stale_after_seconds = env_int("INGESTION_PROCESSING_STALE_SECONDS", 600)
    claimed = await parse_run_repo.claim_run_for_processing(run_id, stale_after_seconds=stale_after_seconds)
    if not claimed:
        logger.bind(run_id=run_id, status=run.status).info("Skipping duplicate ingestion job execution")
        return

    retry_limit = max(0, env_int("INGESTION_RETRY_LIMIT", 2))
    last_error = ""

    for attempt in range(0, retry_limit + 1):
        try:
            attempt_started = perf_counter()
            await ingestion_service.process_ingestion(run_id, request)
            logger.bind(
                run_id=run_id,
                repo_id=request.repo_id,
                doc_path=request.doc_path,
                attempts=attempt + 1,
                retry_limit=retry_limit,
                attempt_duration_ms=round((perf_counter() - attempt_started) * 1000, 2),
                total_duration_ms=round((perf_counter() - started_at) * 1000, 2),
            ).info("Ingestion worker job completed")
            return
        except Exception as exc:
            last_error = str(exc)
            is_retryable = _is_retryable_processing_error(exc)
            if not is_retryable or attempt >= retry_limit:
                logger.bind(
                    run_id=run_id,
                    repo_id=request.repo_id,
                    doc_path=request.doc_path,
                    attempts=attempt + 1,
                    retry_limit=retry_limit,
                    retryable=is_retryable,
                    error_type=type(exc).__name__,
                    total_duration_ms=round((perf_counter() - started_at) * 1000, 2),
                ).error(f"Ingestion job failed permanently: {exc}")
                raise

            next_attempt = attempt + 1
            backoff_seconds = _retry_backoff_seconds(next_attempt)
            next_retry_at = datetime.utcnow() + timedelta(seconds=backoff_seconds)
            await parse_run_repo.mark_run_retrying(
                run_id=run_id,
                retry_count=next_attempt,
                error=last_error,
                next_retry_at=next_retry_at,
            )
            logger.bind(
                run_id=run_id,
                repo_id=request.repo_id,
                doc_path=request.doc_path,
                retry_count=next_attempt,
                retry_limit=retry_limit,
                backoff_seconds=backoff_seconds,
                next_retry_at=next_retry_at.isoformat(),
                error_type=type(exc).__name__,
            ).warning("Retrying ingestion job after transient failure")
            await asyncio.sleep(backoff_seconds)

            reclaimed = await parse_run_repo.claim_run_for_processing(
                run_id,
                stale_after_seconds=stale_after_seconds,
            )
            if not reclaimed:
                logger.bind(run_id=run_id).info("Stopping retries because job is already being processed")
                return
