from typing import Optional

from app.core.database import get_client
from app.core.env import env_bool, env_str
from app.repositories.indexing_run_errors import IndexingRunNotReadyError


def _indexing_runs_enabled() -> bool:
    return env_bool("INDEXING_RUNS_ENABLED", True)


def _indexing_runs_db():
    database = env_str("INDEXING_RUNS_DATABASE", "adpilot_repo_sync")
    return get_client()[database]


async def increment_processed_docs(repo_id: str, snapshot_id: str) -> Optional[dict]:
    if not _indexing_runs_enabled():
        return None

    db = _indexing_runs_db()
    existing = await db.indexing_runs.find_one(
        {"repo_id": repo_id, "snapshot_id": snapshot_id},
        projection={"processed_docs_files": 1, "expected_docs_files": 1},
    )
    if not existing:
        raise IndexingRunNotReadyError(
            f"indexing run not found for repo_id={repo_id} snapshot_id={snapshot_id}"
        )

    update: dict = {"$inc": {"processed_docs_files": 1}}
    if existing.get("processed_docs_files", 0) == 0:
        update["$set"] = {"stages.docs_parse": "running"}

    result = await db.indexing_runs.find_one_and_update(
        {"repo_id": repo_id, "snapshot_id": snapshot_id},
        update,
        return_document=True,
    )
    if not result:
        raise IndexingRunNotReadyError(
            f"indexing run not found for repo_id={repo_id} snapshot_id={snapshot_id}"
        )

    processed = result.get("processed_docs_files", 0)
    expected = result.get("expected_docs_files", 0)
    if expected > 0 and processed >= expected:
        await db.indexing_runs.update_one(
            {"repo_id": repo_id, "snapshot_id": snapshot_id},
            {"$set": {"stages.docs_parse": "completed"}},
        )

    return result
