"""FastAPI Status Listener.

Consumer group: fastlane.status.listener
Topics:
    fastlane.tank.allocate.completed
    fastlane.tank.allocate.failed
    fastlane.tank.allocate.blocked

For each event, POSTs a status update to the ERP callback endpoint so the
ERP UI can render green tick / red cross / blocked.
"""
from __future__ import annotations

import asyncio
import json
import signal

from aiokafka import AIOKafkaConsumer

from app.config import get_settings
from app.logger import get_logger
from app.schemas import ErpStatusUpdate, TankAllocationStatusEvent
from app.services.erp_client import ErpClient

log = get_logger("status-listener")


async def _handle(event_json: dict, erp: ErpClient) -> None:
    event = TankAllocationStatusEvent.model_validate(event_json)
    log.info(
        "Consuming %s plate=%s status=%s fastlane_id=%s",
        event.event_type, event.licence_number, event.status, event.fastlane_id,
    )
    update = ErpStatusUpdate(
        job_id=event.job_id,
        tank_id=event.tank_id,
        licence_number=event.licence_number,
        status=event.status,
        fastlane_id=event.fastlane_id,
        incremental_version=event.incremental_version,
        message=event.message,
        occurred_at=event.occurred_at,
    )
    try:
        await erp.push_status(update)
    except Exception as exc:  # noqa: BLE001
        log.exception("Failed to push status to ERP: %s", exc)
        raise


async def run() -> None:
    s = get_settings()
    topics = [
        s.topic_allocate_completed,
        s.topic_allocate_failed,
        s.topic_allocate_blocked,
    ]
    consumer = AIOKafkaConsumer(
        *topics,
        bootstrap_servers=s.kafka_bootstrap_servers,
        group_id=s.consumer_group_status,
        client_id=f"{s.kafka_client_id}-status",
        enable_auto_commit=False,
        auto_offset_reset="earliest",
        value_deserializer=lambda v: json.loads(v.decode("utf-8")),
    )
    erp = ErpClient()
    await consumer.start()
    log.info(
        "Status listener started. topics=%s group=%s",
        topics, s.consumer_group_status,
    )
    stop = asyncio.Event()

    def _sig(*_):  # noqa: ANN001
        stop.set()

    for sig in (signal.SIGINT, signal.SIGTERM):
        try:
            asyncio.get_event_loop().add_signal_handler(sig, _sig)
        except NotImplementedError:
            pass

    try:
        while not stop.is_set():
            batch = await consumer.getmany(timeout_ms=1000, max_records=20)
            for _tp, msgs in batch.items():
                for msg in msgs:
                    try:
                        await _handle(msg.value, erp)
                    except Exception:  # noqa: BLE001
                        log.exception("Unhandled error processing status event")
                await consumer.commit()
    finally:
        await consumer.stop()
        log.info("Status listener stopped.")


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