"""
第 5 章 - pause / resume 演示
==================================================

场景：模拟下游 DB 写入压力大，对部分分区做背压。

要点：
- pause(partitions) 后这些分区不再返回数据
- 但仍然要持续 poll() 触发心跳，否则会被踢出组
- 等下游缓过来再 resume

业务还有一个高频用法：处理一个长任务（如导出报表）时
pause 全部分区，主线程持续 poll() 保心跳，处理完再 resume。

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

from __future__ import annotations

import os
import signal
import sys
import time

from confluent_kafka import Consumer, KafkaError

BOOTSTRAP = sys.argv[1] if len(sys.argv) > 1 else "127.0.0.1:9092"
TOPIC = "learn.05.orders"
GROUP = f"ch5-pause-resume-{int(time.time())}"

running = True


def stop(*_):
    global running
    print("\n[stop] received signal")
    running = False


signal.signal(signal.SIGINT, stop)
signal.signal(signal.SIGTERM, stop)


class FakeDB:
    """模拟下游 DB：每 5 秒进入一次 30 秒的 backpressure 状态"""

    def __init__(self):
        self._start = time.time()

    def is_busy(self) -> bool:
        elapsed = (time.time() - self._start) % 35
        return 5 < elapsed < 15

    def write(self, msg) -> None:
        time.sleep(0.005)


def main() -> None:
    consumer = Consumer(
        {
            "bootstrap.servers": BOOTSTRAP,
            "group.id": GROUP,
            "auto.offset.reset": "earliest",
            "enable.auto.commit": False,
            "session.timeout.ms": 30000,
            "max.poll.interval.ms": 300000,
        }
    )
    consumer.subscribe([TOPIC])
    db = FakeDB()
    paused = False
    processed = 0

    try:
        while running:
            now_busy = db.is_busy()

            if now_busy and not paused:
                assignment = consumer.assignment()
                if assignment:
                    consumer.pause(assignment)
                    paused = True
                    print(f"[pause] DB busy, paused {len(assignment)} partitions; will keep polling for heartbeat")
            elif not now_busy and paused:
                assignment = consumer.assignment()
                if assignment:
                    consumer.resume(assignment)
                    paused = False
                    print(f"[resume] DB idle, resumed {len(assignment)} partitions")

            msg = consumer.poll(1.0)
            if msg is None:
                if paused:
                    print("  [paused] heartbeat-only poll, no msg returned (expected)")
                continue
            if msg.error():
                if msg.error().code() != KafkaError._PARTITION_EOF:
                    print(f"  err: {msg.error()}")
                continue

            db.write(msg)
            processed += 1
            if processed % 20 == 0:
                consumer.commit(asynchronous=False)
                print(f"  [progress] {processed} msgs processed and committed")
    finally:
        consumer.close()
        print(f"\n[done] processed={processed}")


if __name__ == "__main__":
    main()
