"""
第 4 章 - 幂等 Producer 演示
=================================

演示要点：
1) 开启 enable.idempotence=true 后，Producer 启动会先发
   InitProducerIdRequest 拿到 PID + Epoch（看 librdkafka 日志）。
2) 即使同一条消息被 client 在重试时发出多次，broker 也会按
   (PID, partition, sequence) 做去重，最终 log 里只出现一次。
3) 开了幂等会自动强制：acks=all、max.in.flight <= 5、retries > 0。
   下面的代码故意把 acks=1 注释掉，演示「即使你写了 acks=1 也会被改回 all」。

可以用 kafka-dump-log.sh 验证：
    kafka-dump-log.sh \\
        --files /tmp/kraft-combined-logs/learn.04.idempotent-0/00000000000000000000.log \\
        --print-data-log
看每条消息都有 producerId / sequence 字段。

运行：
    bash ../init.sh
    python idempotent_producer.py
"""

from __future__ import annotations

import sys
import time

from confluent_kafka import Producer

BOOTSTRAP = sys.argv[1] if len(sys.argv) > 1 else "127.0.0.1:9092"
TOPIC = "learn.04.idempotent"
N = 50


def on_delivery(err, msg):
    if err is not None:
        print(f"  FAIL: {err}")
    else:
        print(
            f"  OK  : P{msg.partition()}@{msg.offset()} "
            f"key={msg.key().decode()} val={msg.value().decode()}"
        )


def main() -> None:
    print("[ idempotent producer ] enable.idempotence=true")
    print("librdkafka will enforce: acks=all, max.in.flight<=5, retries>0\n")

    p = Producer(
        {
            "bootstrap.servers": BOOTSTRAP,
            "client.id": "ch4-idempotent-producer",
            "enable.idempotence": True,
            "linger.ms": 5,
            "compression.type": "zstd",
        }
    )

    print(f"Sending {N} messages to {TOPIC}")
    for i in range(N):
        p.produce(
            TOPIC,
            key=f"order-{i % 5}".encode(),
            value=f"payload-#{i}".encode(),
            on_delivery=on_delivery,
        )
        p.poll(0)
        time.sleep(0.01)

    p.flush(15)
    print("\nDone. Inspect the log with kafka-dump-log.sh; each record carries producerId / sequence.")


if __name__ == "__main__":
    main()
