"""
第 3 章 - 最小可用 Producer
==================================

30 行写完一个生产者，演示：
- bootstrap.servers / client.id / acks / enable.idempotence 配置
- 带 Key 发送（同 Key 必去同分区）
- delivery callback 拿到 broker 的 ack 结果
- flush() 确保程序退出前所有消息送达

运行：
    bash ../init.sh                       # 先创建 Topic
    python simple_producer.py             # 发送 10 条到 learn.03.hello
    python simple_producer.py 50 myhost:9092
"""

from __future__ import annotations

import sys
import time

from confluent_kafka import Producer

BOOTSTRAP = sys.argv[2] if len(sys.argv) > 2 else "127.0.0.1:9092"
N = int(sys.argv[1]) if len(sys.argv) > 1 else 10
TOPIC = "learn.03.hello"

producer = Producer(
    {
        "bootstrap.servers": BOOTSTRAP,
        "client.id": "ch3-simple-producer",
        "acks": "all",
        "enable.idempotence": True,
        "linger.ms": 5,
        "compression.type": "zstd",
    }
)


def on_delivery(err, msg):
    if err is not None:
        print(f"❌ 发送失败 key={msg.key()}: {err}")
    else:
        print(
            f"✅ 已写入 {msg.topic()}-{msg.partition()}@offset={msg.offset()} "
            f"key={msg.key().decode() if msg.key() else None} "
            f"value={msg.value().decode()}"
        )


def main() -> None:
    print(f"🚀 向 {TOPIC} 发送 {N} 条消息（acks=all + 幂等）")
    start = time.perf_counter()
    for i in range(N):
        producer.produce(
            topic=TOPIC,
            key=f"user-{i % 3}".encode(),
            value=f"msg-#{i}-ts={int(time.time() * 1000)}".encode(),
            on_delivery=on_delivery,
        )
        # 主动驱动一次回调，避免回调被 buffer 攒住
        producer.poll(0)

    remaining = producer.flush(timeout=15)
    elapsed = time.perf_counter() - start
    print(
        f"\n📊 完成：{N - remaining}/{N} 成功，耗时 {elapsed * 1000:.1f}ms，"
        f"平均 {N / elapsed:.0f} msg/s"
    )
    if remaining > 0:
        print(f"⚠️  仍有 {remaining} 条未发出（flush 超时）")


if __name__ == "__main__":
    main()
