"""Upsert Minion — one job: `newAsset` or `updateAsset` on BASF.

Consumes: fastlane.asset.lookup.done
Produces: fastlane.asset.upserted
DLQ    : fastlane.asset.upsert.dlq
"""
from __future__ import annotations

import asyncio

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

log = get_logger("upsert-minion")
_client = FastlaneClient()


async def handle(event: dict) -> None:
    s = get_settings()
    asset_id = event.get("fastlane_asset_id")
    inc = event.get("incremental_version")
    location_list = event.get("location_list") or FastlaneClient.default_location_list(
        bool(event.get("is_dangerous_goods"))
    )
    owner = event.get("owner") or _client.default_owner()

    try:
        if asset_id is None:
            payload = await _client.new_asset(
                vehicle_type=int(event.get("vehicle_type") or s.fastlane_default_vehicle_type),
                licence_number=event["licence_number"],
                nationality=event.get("nationality"),
                location_list=location_list,
                owner=owner,
                fields=event.get("fields") or None,
            )
            was_created = True
        else:
            payload = await _client.update_asset(
                asset_id=int(asset_id),
                incremental_version=int(inc or 1),
                fields=event.get("fields") or {},
                location_list=location_list,
            )
            was_created = False
    except FastlaneApiError as exc:
        # Let the runner count the attempt / route to DLQ.
        raise RuntimeError(f"BASF upsert failed: {exc}") from exc

    vehicle = (payload or {}).get("vehicle") or {}
    out = dict(event)
    out["event_type"] = s.topic_asset_upserted
    out["fastlane_asset_id"] = vehicle.get("id") or asset_id
    out["incremental_version"] = vehicle.get("incrementalVersion", inc)
    out["was_created"] = was_created
    out["location_list"] = location_list
    log.info(
        "UPSERT %s plate=%s id=%s v=%s",
        "CREATE" if was_created else "UPDATE",
        event.get("licence_number"), out["fastlane_asset_id"], out["incremental_version"],
    )
    await emit(s.topic_asset_upserted, envelope_key(out), out)


async def main() -> None:
    s = get_settings()
    await run_minion(
        name="upsert-minion",
        topic=s.topic_asset_lookup_done,
        group_id=s.consumer_group_upsert,
        handler=handle,
        dlq_topic=s.topic_upsert_dlq,
    )


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