#!/usr/bin/env python3
"""
consume_process_produce.py
==========================
Streams 灵魂模式：从 input topic 读 → 处理 → 写 output topic + 提交 input offset，
全部在一个事务里。这是 Kafka EOS 协议层最重要的使用场景。

用法：
  pip install confluent-kafka

  # 1) 灌一些原始数据
  python consume_process_produce.py seed 200

  # 2) 启动事务式 stream app（持续运行，Ctrl+C 退出）
  python consume_process_produce.py run

  # 3) 用 read_committed_consumer.py 看 output topic
  python read_committed_consumer.py learn.13.cpp_out committed
"""

import os
import sys
import time
import json

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

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
T_IN = "learn.13.cpp_in"
T_OUT = "learn.13.cpp_out"
TX_ID = "tx-cpp-stream-1"


def ensure():
    a = AdminClient({"bootstrap.servers": BOOTSTRAP})
    existing = a.list_topics(timeout=5).topics
    new = [NewTopic(t, num_partitions=3, replication_factor=1)
           for t in (T_IN, T_OUT) if t not in existing]
    if new:
        a.create_topics(new); time.sleep(1)


def seed(n):
    ensure()
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    for i in range(n):
        p.produce(T_IN, key=str(i % 5), value=json.dumps({"raw": i}).encode())
    p.flush(10)
    print(f"seeded {n} raw events to {T_IN}")


def run():
    ensure()
    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "cpp-stream-group",
        "enable.auto.commit": False,
        "isolation.level": "read_committed",
        "auto.offset.reset": "earliest",
    })
    consumer.subscribe([T_IN])

    producer = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "enable.idempotence": True,
        "transactional.id": TX_ID,
        "transaction.timeout.ms": 60000,
        "linger.ms": 5,
    })
    producer.init_transactions()

    print(f"running consume-process-produce: {T_IN} → {T_OUT}, tx_id={TX_ID}")
    try:
        while True:
            msgs = consumer.consume(num_messages=50, timeout=1.0)
            msgs = [m for m in msgs if m and not m.error()]
            if not msgs:
                continue

            producer.begin_transaction()
            try:
                for m in msgs:
                    raw = json.loads(m.value())
                    out = {"clean": raw["raw"] * 2, "src_off": m.offset(), "src_p": m.partition()}
                    producer.produce(T_OUT, key=m.key(), value=json.dumps(out).encode())

                # 把 input offsets 包进同一个事务
                producer.send_offsets_to_transaction(
                    consumer.position(consumer.assignment()),
                    consumer.consumer_group_metadata(),
                )
                producer.commit_transaction()
                print(f"✅ tx commit: processed {len(msgs)} msgs (offsets推进 input topic)")
            except KafkaException as e:
                err = e.args[0]
                if err.txn_requires_abort():
                    producer.abort_transaction()
                    print(f"⚠️ tx aborted: {err}")
                    continue
                else:
                    raise
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    if len(sys.argv) < 2:
        print(__doc__); sys.exit(0)
    if sys.argv[1] == "seed":
        seed(int(sys.argv[2]) if len(sys.argv) > 2 else 100)
    elif sys.argv[1] == "run":
        run()
    else:
        print(__doc__)
