#!/usr/bin/env python3
"""
tombstone_delete.py
===================
演示用 null value 实现「删除一个 Key」的完整流程：
  1) 先写一些 (key, value) 数据
  2) 再写 (key, null) 作为 Tombstone
  3) 等若干秒让 Cleaner 触发
  4) 从头消费，看到 Key 物理消失（或者剩一个 tombstone）

用法：
  pip install confluent-kafka
  bash ../init.sh
  python tombstone_delete.py demo

观察：
  - 第一次 read：Key=item-001 还有数据
  - 写 tombstone 后立刻 read：能看到 (item-001, null) + 历史 value
  - 等 Cleaner 触发：历史 value 消失，只剩 tombstone
  - 等 delete.retention.ms (init.sh 里设 10s) 过期 + Cleaner 再跑：tombstone 也消失
"""

import os
import sys
import time
import json

from confluent_kafka import Producer, Consumer

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = os.environ.get("TOPIC", "learn.14.compact")


def produce(records):
    """records: list of (key_str, value_str | None)"""
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    for k, v in records:
        p.produce(TOPIC,
                  key=k.encode() if k else None,
                  value=v.encode() if v is not None else None)   # ← None 即 Tombstone
    p.flush(10)


def read_snapshot(label):
    print(f"\n--- snapshot: {label} ---")
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "ts-snap-" + str(os.getpid()) + "-" + str(time.time()),
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    c.subscribe([TOPIC])
    items = {}
    deadline = time.time() + 4
    while time.time() < deadline:
        msg = c.poll(1.0)
        if msg is None: continue
        if msg.error(): continue
        deadline = time.time() + 1.5
        k = msg.key().decode() if msg.key() else "<null-key>"
        v = msg.value()
        items.setdefault(k, []).append({
            "offset": msg.offset(),
            "value": v.decode() if v else "TOMBSTONE",
        })
    c.close()
    for k in sorted(items.keys()):
        print(f"  Key={k!r}:")
        for r in items[k]:
            print(f"    offset={r['offset']} value={r['value']}")
    if not items:
        print("  (空)")


def wait(secs):
    print(f"\nsleeping {secs}s ...")
    for i in range(secs, 0, -1):
        sys.stdout.write(f"\r  remaining: {i:3d}s ")
        sys.stdout.flush()
        time.sleep(1)
    sys.stdout.write("\n")


def demo():
    print("STEP 1: 写入 3 个 Key 各 5 个版本")
    rs = []
    for v in range(5):
        for k in ["user-1", "user-2", "user-3"]:
            rs.append((k, json.dumps({"v": v, "name": f"{k}@v{v}"})))
    produce(rs)
    read_snapshot("写入后")

    print("\nSTEP 2: 给 user-2 发 Tombstone")
    produce([("user-2", None)])
    read_snapshot("发完 tombstone 立刻读")

    print("\nSTEP 3: 等 25 秒让 Cleaner 跑（init.sh 里 segment.ms=5s, dirty ratio=0.1）")
    wait(25)
    read_snapshot("Cleaner 跑过后（user-2 应该只剩 tombstone；user-1/3 应该只剩最新版本）")

    print("\nSTEP 4: 再等 25 秒让 delete.retention.ms (10s) 过期 + 下一轮 Cleaner")
    wait(25)
    read_snapshot("最终态（user-2 的 tombstone 也物理消失，只剩 user-1/3 的最新版本）")


def main():
    if len(sys.argv) < 2 or sys.argv[1] == "demo":
        demo()
    elif sys.argv[1] == "snapshot":
        read_snapshot(sys.argv[2] if len(sys.argv) > 2 else "now")
    else:
        print(__doc__)


if __name__ == "__main__":
    main()
