#!/usr/bin/env python3
"""
Schema Evolution 演示：
1) 用 v1 注册 → 写 1 条
2) 用 v2（加 email + age 都带默认值）注册 → 应通过（BACKWARD）
3) 故意提交一个不兼容的 v3（删除 name 字段且无默认值）→ 应被 Schema Registry 拒绝
4) 把兼容性临时改成 NONE，再注册 → 通过（演示「逃生通道」）
5) 改回 BACKWARD

依赖：pip install confluent-kafka[avro] requests
"""

import json
import os
import sys
import requests

from confluent_kafka.schema_registry import SchemaRegistryClient, Schema

SR_URL = os.getenv("SR_URL", "http://localhost:8081")
SUBJECT = "learn.18.users-value"


def get_compat(subject: str = None) -> str:
    """读取当前兼容性策略。"""
    url = f"{SR_URL}/config/{subject}" if subject else f"{SR_URL}/config"
    r = requests.get(url)
    if r.status_code == 404:
        # subject 没有覆盖，回退全局
        return get_compat()
    return r.json().get("compatibilityLevel", "UNKNOWN")


def set_compat(level: str, subject: str = None):
    url = f"{SR_URL}/config/{subject}" if subject else f"{SR_URL}/config"
    r = requests.put(url, json={"compatibility": level},
                     headers={"Content-Type": "application/vnd.schemaregistry.v1+json"})
    r.raise_for_status()
    print(f"  ✓ 兼容性已设为 {level} ({'subject' if subject else 'global'})")


def register(subject: str, schema_dict: dict, label: str):
    sr = SchemaRegistryClient({"url": SR_URL})
    schema_str = json.dumps(schema_dict)
    try:
        sid = sr.register_schema(subject, Schema(schema_str, "AVRO"))
        print(f"  ✓ [{label}] 注册成功，Schema ID = {sid}")
        return sid
    except Exception as e:
        print(f"  ✗ [{label}] 注册失败：{e}")
        return None


def main():
    print(f"Schema Registry URL = {SR_URL}")
    print(f"Subject = {SUBJECT}\n")

    # ----- 0. 清理 -----
    requests.delete(f"{SR_URL}/subjects/{SUBJECT}?permanent=false")
    print("已清理旧版本（软删）\n")

    # ----- 1. 注册 v1 -----
    print("[Step 1] 注册 v1：基础结构")
    v1 = {
        "type": "record", "name": "User",
        "fields": [
            {"name": "id", "type": "long"},
            {"name": "name", "type": "string"},
        ],
    }
    register(SUBJECT, v1, "v1")

    # ----- 2. 当前兼容性 -----
    cur = get_compat(SUBJECT)
    print(f"\n当前兼容性策略：{cur}\n")

    # ----- 3. 注册 v2（兼容） -----
    print("[Step 2] 注册 v2：加 email、age（带默认值）→ BACKWARD 应通过")
    v2 = {
        "type": "record", "name": "User",
        "fields": [
            {"name": "id", "type": "long"},
            {"name": "name", "type": "string"},
            {"name": "email", "type": "string", "default": "unknown@example.com"},
            {"name": "age", "type": ["null", "int"], "default": None},
        ],
    }
    register(SUBJECT, v2, "v2 兼容变更")

    # ----- 4. 注册 v3（不兼容） -----
    print("\n[Step 3] 注册 v3：删除 name 字段（无默认值）→ BACKWARD 应被拒绝")
    v3_bad = {
        "type": "record", "name": "User",
        "fields": [
            {"name": "id", "type": "long"},
            {"name": "email", "type": "string", "default": "unknown@example.com"},
            {"name": "age", "type": ["null", "int"], "default": None},
        ],
    }
    register(SUBJECT, v3_bad, "v3 破坏性变更")

    # ----- 5. 临时改 NONE 再注册 -----
    print("\n[Step 4] 临时把 Subject 兼容性改为 NONE 再注册")
    set_compat("NONE", SUBJECT)
    register(SUBJECT, v3_bad, "v3 NONE 模式")

    # ----- 6. 改回 BACKWARD -----
    print("\n[Step 5] 改回 BACKWARD")
    set_compat("BACKWARD", SUBJECT)

    # ----- 7. 列出所有版本 -----
    versions = requests.get(f"{SR_URL}/subjects/{SUBJECT}/versions").json()
    print(f"\n最终版本列表：{versions}")


if __name__ == "__main__":
    try:
        main()
    except requests.HTTPError as e:
        print(f"HTTP error: {e.response.text}", file=sys.stderr)
        sys.exit(1)
