#!/usr/bin/env python3
"""
inspect_consumer_offsets.py
===========================
直接消费 Kafka 内部 Topic `__consumer_offsets`，解析其二进制 Key/Value 结构，
打印出每条 OffsetCommit / GroupMetadata 记录。

这是「Kafka 自己吃自己的狗粮」最直观的一个演示：
  - Key 描述「这是哪个 group 的哪个 (topic, partition) 的位点」
  - Value 描述「offset / leader_epoch / metadata / commit_timestamp」
  - 整个 topic 用 cleanup.policy=compact 保留每个 Key 的最新 Value

二进制结构（OffsetCommit, version 3+）：
  Key:
    int16 version
    string group_id        (int16 length + bytes)
    string topic
    int32  partition
  Value:
    int16 version
    int64 offset
    int32 leader_epoch     (v3+，无则 -1)
    string metadata
    int64 commit_timestamp
    int64 expire_timestamp (旧版才有)

用法：
  pip install confluent-kafka
  python inspect_consumer_offsets.py

  # 然后另一个终端跑任意消费者并 commit，回到这里看到对应输出。
"""

import os
import struct
import logging

from confluent_kafka import Consumer

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")

logging.basicConfig(level=logging.INFO, format="%(asctime)s %(message)s", datefmt="%H:%M:%S")
log = logging.getLogger("inspect")


class Reader:
    def __init__(self, buf):
        self.buf = buf; self.pos = 0
    def i16(self):
        v = struct.unpack(">h", self.buf[self.pos:self.pos+2])[0]; self.pos += 2; return v
    def i32(self):
        v = struct.unpack(">i", self.buf[self.pos:self.pos+4])[0]; self.pos += 4; return v
    def i64(self):
        v = struct.unpack(">q", self.buf[self.pos:self.pos+8])[0]; self.pos += 8; return v
    def str(self):
        n = self.i16()
        if n < 0: return None
        v = self.buf[self.pos:self.pos+n].decode("utf-8", errors="replace"); self.pos += n; return v
    def remain(self): return len(self.buf) - self.pos


def parse_offset_key(key):
    r = Reader(key)
    version = r.i16()
    # version 0 / 1 = OffsetCommit
    # version 2     = GroupMetadata
    if version not in (0, 1):
        return None
    group = r.str()
    topic = r.str()
    partition = r.i32()
    return {"kind": "OffsetCommit", "version": version,
            "group": group, "topic": topic, "partition": partition}


def parse_group_meta_key(key):
    r = Reader(key)
    version = r.i16()
    if version != 2:
        return None
    group = r.str()
    return {"kind": "GroupMetadata", "version": version, "group": group}


def parse_offset_value(val):
    if val is None:
        return {"tombstone": True}
    r = Reader(val)
    try:
        version = r.i16()
        offset = r.i64()
        leader_epoch = r.i32() if version >= 3 else -1
        metadata = r.str()
        commit_ts = r.i64()
        expire_ts = r.i64() if version <= 1 else None
        return {"version": version, "offset": offset, "leader_epoch": leader_epoch,
                "metadata": metadata, "commit_timestamp": commit_ts,
                "expire_timestamp": expire_ts}
    except Exception as e:
        return {"raw_len": len(val), "parse_error": str(e)}


def main():
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "inspect-internal-" + str(os.getpid()),
        "enable.auto.commit": False,
        "auto.offset.reset": "latest",        # 只看新提交，避免历史刷屏
    })
    c.subscribe(["__consumer_offsets"])

    log.info("listening to __consumer_offsets ... 在另一个终端跑任意 commit 即可看到输出")
    try:
        while True:
            msg = c.poll(1.0)
            if msg is None: continue
            if msg.error():
                log.error(msg.error()); continue

            key = msg.key() or b""
            val = msg.value()
            if len(key) < 2:
                continue

            # 先尝试当作 OffsetCommit
            parsed_key = parse_offset_key(key)
            if parsed_key:
                parsed_val = parse_offset_value(val)
                log.info(
                    f"📌 OffsetCommit  partition_in_internal=p{msg.partition()}  "
                    f"group={parsed_key['group']!r}  topic={parsed_key['topic']!r}  "
                    f"part={parsed_key['partition']}  →  offset={parsed_val.get('offset')}  "
                    f"leader_epoch={parsed_val.get('leader_epoch')}  "
                    f"commit_ts={parsed_val.get('commit_timestamp')}"
                    f"{'  [TOMBSTONE]' if parsed_val.get('tombstone') else ''}"
                )
                continue

            # 否则尝试 GroupMetadata
            gm = parse_group_meta_key(key)
            if gm:
                size = len(val) if val else 0
                log.info(
                    f"📁 GroupMetadata  partition_in_internal=p{msg.partition()}  "
                    f"group={gm['group']!r}  value_size={size} bytes  "
                    f"{'(tombstone)' if val is None else '(成员/分配方案 二进制)'}"
                )
                continue

            log.info(f"unknown record key_version={struct.unpack('>h', key[:2])[0]}")
    except KeyboardInterrupt:
        pass
    finally:
        c.close()


if __name__ == "__main__":
    main()
