"""Kafka producer wrapper (async)."""
from __future__ import annotations

import asyncio
import json
from typing import Optional

from aiokafka import AIOKafkaProducer

from .config import get_settings
from .logger import get_logger

log = get_logger(__name__)


class KafkaProducer:
    _instance: Optional["KafkaProducer"] = None
    _lock = asyncio.Lock()

    def __init__(self) -> None:
        s = get_settings()
        self._producer = AIOKafkaProducer(
            bootstrap_servers=s.kafka_bootstrap_servers,
            client_id=s.kafka_client_id,
            # NOTE: idempotence requires a working transaction coordinator on
            # the broker (transaction.state.log.replication.factor <= brokers,
            # transaction.state.log.min.isr <= brokers). Dev/KRaft clusters
            # often have this misconfigured → InitProducerId hangs.
            # Similarly, acks='all' + a dev cluster with stale ISR causes
            # publish to time out. Keep this at 'acks=1' for dev; use 'all'
            # in prod once broker replication factors are correctly tuned.
            acks=1,
            enable_idempotence=False,
            value_serializer=lambda v: json.dumps(v, default=str).encode("utf-8"),
            key_serializer=lambda k: (k or "").encode("utf-8"),
        )
        self._started = False

    async def start(self) -> None:
        if not self._started:
            await self._producer.start()
            self._started = True
            log.info("Kafka producer started")

    async def stop(self) -> None:
        if self._started:
            await self._producer.stop()
            self._started = False
            log.info("Kafka producer stopped")

    async def publish(self, topic: str, key: str, value: dict) -> None:
        await self.start()
        log.info("→ Kafka publish topic=%s key=%s", topic, key)
        await self._producer.send_and_wait(topic, value=value, key=key)


async def get_producer() -> KafkaProducer:
    if KafkaProducer._instance is None:
        async with KafkaProducer._lock:
            if KafkaProducer._instance is None:
                KafkaProducer._instance = KafkaProducer()
                await KafkaProducer._instance.start()
    return KafkaProducer._instance
