"""
Ch7 配套代码 4 / 4 —— Stream + Consumer Group 完整示例

演示：
  1. 生产者持续 XADD 订单消息
  2. 同一 Group 内 2 个 worker 协作分摊消费
  3. ACK 机制 + PEL（待确认列表）观察
  4. 模拟 worker-2 崩溃，未 ACK 的消息由 worker-1 通过 XCLAIM 接管
"""

import threading
import time
import random
import redis

POOL = redis.ConnectionPool(host="127.0.0.1", port=6379, decode_responses=True)

STREAM = "orders"
GROUP = "payproc"
N_MESSAGES = 20


def section(title: str) -> None:
    print("\n" + "=" * 62)
    print(title)
    print("=" * 62)


def setup() -> None:
    r = redis.Redis(connection_pool=POOL)
    r.delete(STREAM)
    try:
        r.xgroup_create(STREAM, GROUP, id="$", mkstream=True)
    except redis.ResponseError as e:
        if "BUSYGROUP" not in str(e): raise
    print(f"  ✅ 流 {STREAM} 与消费组 {GROUP} 已就绪")


def producer() -> None:
    r = redis.Redis(connection_pool=POOL)
    for i in range(N_MESSAGES):
        msg_id = r.xadd(STREAM, {
            "order_id": str(1000 + i),
            "user": f"u{random.randint(1,99):02d}",
            "amount": str(random.randint(10, 999)),
        })
        print(f"  [P] XADD {msg_id}  order_id={1000+i}")
        time.sleep(0.1)
    print("  [P] ✅ 生产完毕")


def worker(name: str, fail_rate: float = 0.0, stop_after: int = -1) -> None:
    """
    fail_rate: 模拟「领走但不 ACK」的概率
    stop_after: 处理 N 条后退出（模拟崩溃）。-1 = 不退出
    """
    r = redis.Redis(connection_pool=POOL)
    handled = 0
    while True:
        msgs = r.xreadgroup(
            groupname=GROUP, consumername=name,
            streams={STREAM: ">"}, count=2, block=2000,
        )
        if not msgs:
            print(f"  [{name}] 无新消息 5s，退出")
            return
        for _stream, entries in msgs:
            for msg_id, body in entries:
                handled += 1
                will_fail = random.random() < fail_rate
                tag = "💥 故意不 ACK" if will_fail else "✅ ACK"
                print(f"  [{name}] got {msg_id}  order={body.get('order_id')}  {tag}")
                time.sleep(random.uniform(0.05, 0.2))
                if not will_fail:
                    r.xack(STREAM, GROUP, msg_id)
                if stop_after > 0 and handled >= stop_after:
                    print(f"  [{name}] 🛑 模拟崩溃，已处理 {handled} 条退出")
                    return


def show_pending() -> None:
    section("Demo: XPENDING —— 谁还有没 ACK 的消息？")
    r = redis.Redis(connection_pool=POOL)
    summary = r.xpending(STREAM, GROUP)
    print(f"  pending 总数 = {summary['pending']}")
    if summary["pending"] == 0:
        print("  ✅ 所有消息都已 ACK")
        return
    print(f"  消费者明细 = {summary['consumers']}")

    detail = r.xpending_range(STREAM, GROUP, "-", "+", count=20)
    for d in detail:
        print(f"    {d['message_id']}  consumer={d['consumer']}  idle={d['time_since_delivered']}ms")
    return detail


def demo_claim(detail) -> None:
    """把 worker-2 的死消息转给 worker-1 重新处理"""
    if not detail: return
    section("Demo: XCLAIM —— 把死信转给 worker-1 重新处理")
    r = redis.Redis(connection_pool=POOL)

    dead_ids = [d["message_id"] for d in detail if d["consumer"] != "worker-1"]
    if not dead_ids:
        print("  没有需要 claim 的消息")
        return

    print(f"  XCLAIM {len(dead_ids)} 条 → worker-1（min-idle=0 强制接管）")
    claimed = r.xclaim(STREAM, GROUP, "worker-1", min_idle_time=0, message_ids=dead_ids)
    for msg_id, body in claimed:
        print(f"    worker-1 接管 {msg_id}  order={body.get('order_id')}  → ACK")
        r.xack(STREAM, GROUP, msg_id)


def main() -> None:
    section("Stream Consumer Group 完整流程")
    setup()

    section("Demo: 1 个 Producer + 2 个 Worker 协作消费")
    threads = [
        threading.Thread(target=producer, name="P"),
        threading.Thread(target=worker, args=("worker-1", 0.0, -1), name="W1"),
        threading.Thread(target=worker, args=("worker-2", 0.4, 5), name="W2"),
    ]
    for t in threads: t.start()
    for t in threads: t.join()

    detail = show_pending()
    demo_claim(detail)
    show_pending()

    section("收尾")
    r = redis.Redis(connection_pool=POOL)
    info = r.xinfo_stream(STREAM)
    print(f"  Stream length     = {info['length']}")
    print(f"  Last generated id = {info['last-generated-id']}")
    print(f"  ✅ Demo 完成")
    r.delete(STREAM)


if __name__ == "__main__":
    try:
        main()
    except redis.ConnectionError as e:
        print(f"❌ Redis 连接失败: {e}")
