#!/usr/bin/env python3
"""
min_isr_demo.py — 演示 min.insync.replicas 起作用的全过程

剧本：
    1) 创建 Topic：3 副本，min.insync.replicas=2；
    2) Producer 配置 acks=all；
    3) 正常情况下 ISR=3，写入成功；
    4) 用 docker compose stop 停掉两个 Broker，让 ISR 缩到 1；
    5) 此时 Producer 写入会拿到 NotEnoughReplicas / NotEnoughReplicasAfterAppend 错误；
    6) 把 Broker 起回来，ISR 恢复到 ≥ 2，写入恢复成功；
    7) 输出全过程的 ack / 错误统计。

用法：
    python min_isr_demo.py
    python min_isr_demo.py --bootstrap localhost:9092 --topic learn.09.minisr
"""

from __future__ import annotations

import argparse
import subprocess
import time
from datetime import datetime

from confluent_kafka import KafkaException, Producer
from confluent_kafka.admin import AdminClient, NewTopic, ConfigResource, ResourceType


def parse_args():
    p = argparse.ArgumentParser()
    p.add_argument("--bootstrap", default="localhost:9092,localhost:9094,localhost:9096")
    p.add_argument("--topic", default="learn.09.minisr")
    p.add_argument("--compose-dir", default="/data/workspace/learnNote/kafka")
    return p.parse_args()


def ensure_topic(admin: AdminClient, topic: str):
    md = admin.list_topics(timeout=10)
    if topic in md.topics and md.topics[topic].error is None:
        print(f"[init] topic {topic} 已存在")
        return
    nt = NewTopic(topic, num_partitions=1, replication_factor=3,
                  config={"min.insync.replicas": "2"})
    fs = admin.create_topics([nt])
    for t, f in fs.items():
        f.result()
        print(f"[init] 创建 topic {t} (1 part, RF=3, min.isr=2)")


def show_isr(admin, topic):
    md = admin.list_topics(topic=topic, timeout=10)
    p = md.topics[topic].partitions[0]
    print(f"[meta] {topic}-P0  Leader={p.leader}  ISR={list(p.isrs)}  "
          f"Replicas={list(p.replicas)}")
    return p


def docker_compose(cmd, cwd):
    full = ["docker", "compose"] + cmd
    print(f"[shell] cd {cwd} && {' '.join(full)}")
    return subprocess.run(full, cwd=cwd, check=False,
                          stdout=subprocess.PIPE, stderr=subprocess.STDOUT)


class Counter:
    def __init__(self):
        self.ok = 0
        self.fail = 0
        self.errors = {}

    def on_delivery(self, err, msg):
        if err is None:
            self.ok += 1
        else:
            self.fail += 1
            key = err.name() if hasattr(err, "name") else str(err)
            self.errors[key] = self.errors.get(key, 0) + 1

    def summary(self, tag):
        print(f"[stat:{tag}] 成功={self.ok}  失败={self.fail}  错误明细={self.errors}")


def burst(producer: Producer, topic, n, c: Counter):
    for i in range(n):
        try:
            producer.produce(topic, value=f"{datetime.now().isoformat()}-{i}".encode(),
                             on_delivery=c.on_delivery)
            producer.poll(0)
        except BufferError:
            producer.poll(0.5)
        except KafkaException as e:
            c.fail += 1
            c.errors[str(e)] = c.errors.get(str(e), 0) + 1
    producer.flush(15)


def main():
    args = parse_args()
    admin = AdminClient({"bootstrap.servers": args.bootstrap})

    ensure_topic(admin, args.topic)
    p = show_isr(admin, args.topic)

    producer = Producer({
        "bootstrap.servers": args.bootstrap,
        "acks": "all",
        "enable.idempotence": True,
        "delivery.timeout.ms": 15000,
        "retries": 3,                       # 故意调小，观察失败更明显
        "request.timeout.ms": 5000,
    })

    print("\n=== 阶段 1：稳态写入（ISR=3，min.isr=2 满足） ===")
    c1 = Counter()
    burst(producer, args.topic, 50, c1)
    c1.summary("steady")
    show_isr(admin, args.topic)

    # 停 2 个 follower，让 ISR 缩到 1（只剩 Leader）
    leader_id = p.leader
    others = [b for b in [1, 2, 3] if b != leader_id]
    container_map = {1: "kafka1", 2: "kafka2", 3: "kafka3"}
    targets = [container_map[i] for i in others]

    print(f"\n=== 阶段 2：停掉两个 Follower 容器 {targets}，让 ISR 缩到 1 ===")
    for c in targets:
        docker_compose(["stop", c], args.compose_dir)

    # 等 Controller 把它们踢出 ISR（默认 30s）
    print("[wait] 等待 35s 让 ISR 收缩到 1…")
    for i in range(35):
        time.sleep(1)
        if i % 5 == 0:
            show_isr(admin, args.topic)

    print("\n=== 阶段 3：再写入，应当大量收到 NotEnoughReplicas 错误 ===")
    c2 = Counter()
    burst(producer, args.topic, 50, c2)
    c2.summary("after-shrink")
    show_isr(admin, args.topic)

    print(f"\n=== 阶段 4：把 Broker 起回来，等 ISR 扩张回去 ===")
    for c in targets:
        docker_compose(["start", c], args.compose_dir)
    print("[wait] 等待 60s 让 ISR 重新扩张…")
    for i in range(60):
        time.sleep(1)
        if i % 10 == 0:
            show_isr(admin, args.topic)

    print("\n=== 阶段 5：再写入，应当全部成功 ===")
    c3 = Counter()
    burst(producer, args.topic, 50, c3)
    c3.summary("recovered")
    show_isr(admin, args.topic)

    print("\n=== 总结 ===")
    print(f"  正常阶段 OK={c1.ok} FAIL={c1.fail}")
    print(f"  ISR<min  OK={c2.ok} FAIL={c2.fail}    ← 这里应当大量 NotEnoughReplicas")
    print(f"  恢复阶段 OK={c3.ok} FAIL={c3.fail}")
    print("\n💡 结论：min.insync.replicas 是「acks=all 真正能保住数据」的兜底，"
          "ISR 不够时宁可拒绝写入，也不让数据有丢的风险。")


if __name__ == "__main__":
    main()
