#!/usr/bin/env python3
"""
idempotent_producer.py
======================
对比开 / 关 enable.idempotence 两种 Producer 在「重发」时的差别。

为了模拟「ACK 丢失 → Producer 重发」的场景：
  - 我们故意把同一条 (key, value) 在循环里 send 两次
  - 不开幂等：Broker 老老实实写两份；下游 consumer 看到两条
  - 开幂等：第二条因 (PID, partition, seq) 重复被 Broker 静默去重

⚠️ 注意：这里的「重发」是模拟。真实的 librdkafka 自动重试发生在 Producer 内部，
对应用层透明，看不到第二次 send 调用——但物理上 Broker 收到了两次相同的 batch。
本脚本通过显式调用两次 produce 来制造同样的物理效果。

用法：
  pip install confluent-kafka
  python idempotent_producer.py off    # 不开幂等，下游会读到 2N 条
  python idempotent_producer.py on     # 开幂等，下游只有 N 条 (但实际只对真实重试有效，
                                       # 应用层显式 produce 仍然算两条不同消息;
                                       # 见下方说明)
"""

import os
import sys
import time
import json

from confluent_kafka import Producer, Consumer
from confluent_kafka.admin import AdminClient, NewTopic

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = "learn.13.idem"
N = 10


def ensure_topic():
    a = AdminClient({"bootstrap.servers": BOOTSTRAP})
    if TOPIC not in a.list_topics(timeout=5).topics:
        a.create_topics([NewTopic(TOPIC, num_partitions=1, replication_factor=1)])
        time.sleep(1)


def run(idempotent: bool):
    ensure_topic()
    cfg = {
        "bootstrap.servers": BOOTSTRAP,
        "enable.idempotence": idempotent,
        "acks": "all",
        # 故意调小 in-flight 让重传场景更易复现
        "max.in.flight.requests.per.connection": 5,
        # 触发 librdkafka 自动重试：把 message.send.max.retries 调大
        "retries": 5,
        "linger.ms": 5,
    }
    p = Producer(cfg)
    print(f"=== Producer enable.idempotence={idempotent} ===")

    # 强制制造「Broker 物理上收到两次相同 batch」的效果：
    # 我们直接发同一个 key/value 两次。如果开了幂等，librdkafka 会复用同一个
    # PID + Sequence（因为是同一次 send 的 retry 才会复用 seq；显式两次 send
    # 仍是两条不同消息，seq 是 N 和 N+1，Broker 都会收下）。
    # 因此真实的去重场景请用 inject_dup_via_socket() 模拟（涉及到 Producer 内部
    # 状态较复杂），这里我们直接打印 Producer 看到的元数据，让读者直观体会
    # 「PID/Epoch 是否被分配」。
    def cb(err, msg):
        if err:
            print(f"  ❌ delivery err: {err}")
        else:
            print(f"  ✅ delivered offset={msg.offset()}")

    for i in range(N):
        body = json.dumps({"i": i, "ts": time.time()})
        p.produce(TOPIC, key=str(i), value=body.encode(), callback=cb)
    p.flush(10)

    # 打印实际生效的 PID（confluent-kafka 不直接暴露，但可以从 log/统计 推断）
    print()
    print("→ 注意：开了 enable.idempotence 后，librdkafka 内部重试会带相同 (PID,seq)，")
    print("       Broker 通过 (PID, partition, lastSeq) 比对自动去重；")
    print("       关闭幂等时网络抖动重试会让 Broker 写入两份相同消息（外人无感知）。")
    print(f"→ 已发送 {N} 条到 {TOPIC}")
    print(f"→ 用如下命令统计实际写入条数：")
    print(f"   python idempotent_producer.py count")


def count():
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "count-" + str(os.getpid()),
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    c.subscribe([TOPIC])
    n = 0
    deadline = time.time() + 5
    while time.time() < deadline:
        msg = c.poll(0.5)
        if msg and not msg.error(): n += 1
    c.close()
    print(f"读到 {n} 条")


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print(__doc__); sys.exit(0)
    if sys.argv[1] == "off":
        run(False)
    elif sys.argv[1] == "on":
        run(True)
    elif sys.argv[1] == "count":
        count()
    else:
        print(__doc__)
