"""File-Upload Minion — one job: `uploadFile` on BASF (spec §2.3).

Consumes: fastlane.asset.upserted
Produces: fastlane.asset.docs.uploaded
          fastlane.tank.allocate.completed  (PENDING — awaiting BASF review)
DLQ     : fastlane.asset.docs.dlq

Spec rules enforced here (§2.3 / §2.4):
  - fields[] must have ≥1 slug per document  (§2.3)
  - Each uploadFile call increases incrementalVersion by 1  (§2.3)
  - Allowed types: PDF, JPG, PNG  (§2.3)
  - If docs[] supplied but ALL uploads fail → raise so runner sends to DLQ.
    commitAsset will be rejected by BASF if mandatory fields have no doc. (§2.4)
  - If docs[] is empty → skip upload, proceed to docs.uploaded so commit
    can still run (BASF allows commit when NO field values are set).
"""
from __future__ import annotations

import asyncio

from app.config import get_settings
from app.logger import get_logger
from app.schemas import AllocationStatus, TankAllocationStatusEvent
from app.services.fastlane_client import FastlaneApiError, FastlaneClient
from minions._runner import emit, envelope_key, run_minion

log = get_logger("file-upload-minion")
_client = FastlaneClient()


async def handle(event: dict) -> None:
    s        = get_settings()
    asset_id = int(event["fastlane_asset_id"])
    inc      = int(event.get("incremental_version") or 1)
    docs     = event.get("documents") or []

    # ── No documents to upload ────────────────────────────────────────────────
    # §2.4: commit only fails if a field WITH a value has no validating doc.
    # If the caller sent no docs (e.g. first-time newAsset with no fields),
    # skip upload entirely and let commit proceed.
    if not docs:
        log.info(
            "No documents supplied for asset=%s — skipping uploadFile. "
            "commitAsset will succeed if no field values require validation.",
            asset_id,
        )
        out               = dict(event)
        out["event_type"] = s.topic_asset_docs_uploaded
        out["documents_uploaded"] = 0
        await emit(s.topic_asset_docs_uploaded, envelope_key(out), out)
        _emit_pending(s, event, asset_id, inc, uploaded=0, total=0)
        return

    # ── Upload each document ──────────────────────────────────────────────────
    uploaded = 0
    failed   = 0

    for doc in docs:
        name    = str(doc.get("name") or "document.pdf")
        fields  = list(doc.get("fields") or [])
        content = str(doc.get("content_base64") or "")

        # §2.3: every document must validate ≥1 field
        if not fields:
            log.warning(
                "uploadFile SKIP asset=%s name=%s — fields[] is empty (§2.3 requires ≥1 field)",
                asset_id, name,
            )
            failed += 1
            continue

        try:
            payload = await _client.upload_file(
                asset_id            = asset_id,
                incremental_version = inc,
                fields              = fields,
                name                = name,
                content_base64      = content,
            )
            # §2.3: each uploadFile increases incrementalVersion by 1
            vehicle = (payload or {}).get("vehicle") or {}
            inc     = int(vehicle.get("incrementalVersion", inc))
            uploaded += 1
            log.info(
                "uploadFile OK asset=%s name=%s fields=%s new_v=%s",
                asset_id, name, fields, inc,
            )
        except FastlaneApiError as exc:
            failed += 1
            log.warning(
                "uploadFile FAIL asset=%s name=%s fields=%s: [%s] %s",
                asset_id, name, fields, exc.code, exc.message,
            )

    # ── §2.4 guard: block commit if all docs failed ───────────────────────────
    # commitAsset requires every field-with-value to be validated by a doc.
    # If nothing uploaded, raising here sends the event to DLQ for retry.
    if uploaded == 0:
        raise RuntimeError(
            f"All {failed} document upload(s) failed for asset={asset_id}. "
            "Blocking commitAsset — spec §2.4: all fields with information "
            "must be validated with a document before commit."
        )

    if failed:
        log.warning(
            "asset=%s partial upload: %d uploaded, %d failed. "
            "commitAsset may be rejected by BASF if failed docs cover required fields.",
            asset_id, uploaded, failed,
        )

    # ── Emit docs.uploaded → commit-minion fires after BASF webhook ───────────
    out                       = dict(event)
    out["event_type"]         = s.topic_asset_docs_uploaded
    out["incremental_version"]= inc
    out["documents_uploaded"] = uploaded
    out["documents_failed"]   = failed
    await emit(s.topic_asset_docs_uploaded, envelope_key(out), out)

    _emit_pending(s, event, asset_id, inc, uploaded=uploaded, total=len(docs))


def _emit_pending(s, event: dict, asset_id: int, inc: int,
                  uploaded: int, total: int) -> None:
    """Fire ERP-facing PENDING status so UI renders amber clock."""
    status_event = TankAllocationStatusEvent(
        event_type          = s.topic_allocate_completed,
        job_id              = event.get("job_id"),
        tank_id             = event.get("tank_id"),
        licence_number      = event.get("licence_number"),
        status              = AllocationStatus.PENDING,
        fastlane_id         = asset_id,
        incremental_version = inc,
        message             = (
            f"{uploaded}/{total} doc(s) uploaded — awaiting BASF document review"
            if total else "No docs — awaiting BASF review"
        ),
    )
    asyncio.ensure_future(
        emit(
            s.topic_allocate_completed,
            envelope_key(event),
            status_event.model_dump(mode="json"),
        )
    )


async def main() -> None:
    s = get_settings()
    await run_minion(
        name="file-upload-minion",
        topic=s.topic_asset_upserted,
        group_id=s.consumer_group_file_upload,
        handler=handle,
        dlq_topic=s.topic_file_upload_dlq,
    )


if __name__ == "__main__":
    asyncio.run(main())


async def main() -> None:
    s = get_settings()
    await run_minion(
        name="file-upload-minion",
        topic=s.topic_asset_upserted,
        group_id=s.consumer_group_file_upload,
        handler=handle,
        dlq_topic=s.topic_file_upload_dlq,
    )


if __name__ == "__main__":
    asyncio.run(main())
