"""
Ch14 配套代码 5 / 7 —— 异步下单消费者
================================================================

用途：
  从 Stream `seckill:orders` 消费秒杀消息，模拟「写订单到 MySQL」。
  - 用消费组保证「一个消息只被一个消费者处理」
  - 处理完才 XACK，崩溃重启后自动重做未 ACK 的消息

运行：
  python 05_consumer.py
  Ctrl+C 退出
"""

import sys
import time
import signal
import redis

STREAM_KEY = "seckill:orders"
GROUP      = "order_grp"
CONSUMER   = "consumer-1"

stop = False


def init_group(r: redis.Redis) -> None:
    try:
        r.xgroup_create(STREAM_KEY, GROUP, id="0", mkstream=True)
        print(f"✅ 创建消费组 {GROUP}")
    except redis.ResponseError as e:
        if "BUSYGROUP" in str(e):
            print(f"ℹ️  消费组 {GROUP} 已存在")
        else:
            raise


def write_to_mysql(order_id: str, user_id: str, item_id: str) -> None:
    """模拟写 MySQL：实际项目这里是 INSERT INTO orders ..."""
    time.sleep(0.001)


def send_sms(user_id: str, item_id: str) -> None:
    """模拟发短信"""
    time.sleep(0.0005)


def handle_signal(signum, frame):
    global stop
    stop = True
    print("\n收到退出信号，优雅退出中…")


def main():
    signal.signal(signal.SIGINT, handle_signal)
    signal.signal(signal.SIGTERM, handle_signal)

    r = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
    r.ping()
    init_group(r)

    handled = 0
    failed = 0
    print("=" * 56)
    print(f"  消费者 {CONSUMER} 开始消费 (Ctrl+C 退出)")
    print("=" * 56)

    while not stop:
        try:
            msgs = r.xreadgroup(GROUP, CONSUMER,
                                {STREAM_KEY: ">"},
                                count=200, block=2000)
        except redis.ConnectionError:
            time.sleep(1)
            continue

        if not msgs:
            print(f"  …等待新消息  (已处理 {handled})")
            continue

        _, entries = msgs[0]
        for mid, data in entries:
            try:
                write_to_mysql(data["order_id"], data["user_id"], data["item_id"])
                send_sms(data["user_id"], data["item_id"])
                r.xack(STREAM_KEY, GROUP, mid)
                handled += 1
                if handled % 100 == 0:
                    print(f"  ✓ 已处理 {handled} 单, 最新: order={data['order_id'][:8]} user={data['user_id']}")
            except Exception as e:
                failed += 1
                print(f"  ❌ 处理失败 {mid}: {e}（不 ACK，下次重试）")

    pending = r.xpending(STREAM_KEY, GROUP).get("pending", 0) if False else 0
    print(f"\n退出。处理成功 {handled} 单，失败 {failed} 单")


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