#!/usr/bin/env python3
"""
Order Service：FastAPI + Kafka Producer

职责：
1. 接收 HTTP POST /orders
2. 同步入 MySQL（事务）
3. 异步发 Kafka orders.events（Avro，带 idempotence）

启动：
    uvicorn order_producer:app --host 0.0.0.0 --port 8080 --workers 4

依赖：
    pip install fastapi uvicorn pymysql confluent-kafka[avro]
"""

import os
import json
import time
import uuid
from contextlib import asynccontextmanager

import pymysql
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field
from typing import List, Optional

from confluent_kafka import Producer
from confluent_kafka.serialization import SerializationContext, MessageField, StringSerializer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"

MY_HOST = os.getenv("MY_HOST", "127.0.0.1")
MY_PORT = int(os.getenv("MY_PORT", "3306"))
MY_USER = os.getenv("MY_USER", "root")
MY_PWD = os.getenv("MY_PWD", "root")
MY_DB = "shop"

ORDER_EVENT_SCHEMA = """
{
  "type": "record",
  "namespace": "com.shop.events",
  "name": "OrderEvent",
  "fields": [
    {"name": "order_id",   "type": "string"},
    {"name": "user_id",    "type": "long"},
    {"name": "amount",     "type": "double"},
    {"name": "currency",   "type": "string", "default": "CNY"},
    {"name": "status",     "type": {"type":"enum","name":"OrderStatus",
              "symbols":["CREATED","PAID","CANCELED","REFUNDED"]}},
    {"name": "city",       "type": ["null","string"], "default": null},
    {"name": "created_at", "type": {"type":"long","logicalType":"timestamp-millis"}},
    {"name": "items", "type": {"type":"array","items":{
        "type":"record","name":"Item",
        "fields":[
          {"name":"sku","type":"string"},
          {"name":"qty","type":"int"},
          {"name":"price","type":"double"}
        ]}}}
  ]
}
"""


class Item(BaseModel):
    sku: str
    qty: int = Field(ge=1)
    price: float = Field(ge=0)


class CreateOrder(BaseModel):
    user_id: int
    items: List[Item]
    city: Optional[str] = None


# 全局对象
producer: Producer = None
avro_serializer: AvroSerializer = None
key_serializer = StringSerializer("utf_8")
db_pool = None


def get_conn():
    return pymysql.connect(
        host=MY_HOST, port=MY_PORT, user=MY_USER, password=MY_PWD,
        database=MY_DB, autocommit=False, charset="utf8mb4",
    )


@asynccontextmanager
async def lifespan(app: FastAPI):
    global producer, avro_serializer
    sr = SchemaRegistryClient({"url": SR_URL})
    avro_serializer = AvroSerializer(sr, ORDER_EVENT_SCHEMA, lambda x, ctx: x)
    producer = Producer({
        "bootstrap.servers": BOOTSTRAP,
        "acks": "all",
        "enable.idempotence": True,
        "max.in.flight.requests.per.connection": 5,
        "linger.ms": 10,
        "compression.type": "zstd",
        "delivery.timeout.ms": 120000,
        "client.id": f"order-svc-{uuid.uuid4().hex[:8]}",
    })
    yield
    producer.flush(10)


app = FastAPI(title="Order Service", lifespan=lifespan)


def _delivery_cb(err, msg):
    if err is not None:
        # 实际生产应记录到本地 disk / outbox 重试
        print(f"  [!!!] kafka send failed: {err}")
    else:
        print(f"  [✓] kafka offset = {msg.offset()} partition = {msg.partition()}")


@app.get("/health")
def health():
    return {"status": "ok"}


@app.post("/orders")
def create_order(o: CreateOrder):
    order_id = "ORD-" + uuid.uuid4().hex[:16].upper()
    amount = round(sum(i.qty * i.price for i in o.items), 2)
    items_json = json.dumps([i.model_dump() for i in o.items])

    # 1. 入 MySQL（事务）
    conn = get_conn()
    try:
        with conn.cursor() as cur:
            cur.execute(
                "INSERT INTO orders (order_id, user_id, amount, status, city, items_json) "
                "VALUES (%s, %s, %s, %s, %s, %s)",
                (order_id, o.user_id, amount, "CREATED", o.city, items_json),
            )
            # Outbox 模式（可选）：把事件先写出箱表
            cur.execute(
                "INSERT INTO outbox (aggregate_id, event_type, payload) VALUES (%s, %s, %s)",
                (order_id, "OrderCreated", json.dumps({
                    "order_id": order_id, "user_id": o.user_id, "amount": amount,
                    "city": o.city,
                })),
            )
        conn.commit()
    except Exception as e:
        conn.rollback()
        raise HTTPException(500, f"db error: {e}")
    finally:
        conn.close()

    # 2. 异步发 Kafka（这里直接发；生产可由 Outbox Relay 单独进程发）
    event = {
        "order_id": order_id,
        "user_id": o.user_id,
        "amount": amount,
        "currency": "CNY",
        "status": "CREATED",
        "city": o.city,
        "created_at": int(time.time() * 1000),
        "items": [i.model_dump() for i in o.items],
    }
    try:
        value = avro_serializer(event, SerializationContext(TOPIC, MessageField.VALUE))
        producer.produce(
            topic=TOPIC,
            key=key_serializer(order_id, SerializationContext(TOPIC, MessageField.KEY)),
            value=value,
            on_delivery=_delivery_cb,
        )
        producer.poll(0)
    except Exception as e:
        # 发送失败不能让用户感知（已入库）；记日志，由 Outbox Relay 兜底
        print(f"[WARN] enqueue kafka fail (will be relayed by outbox): {e}")

    return {"order_id": order_id, "status": "CREATED", "amount": amount}


@app.get("/orders/{order_id}")
def get_order(order_id: str):
    conn = get_conn()
    try:
        with conn.cursor(pymysql.cursors.DictCursor) as cur:
            cur.execute("SELECT * FROM orders WHERE order_id=%s", (order_id,))
            row = cur.fetchone()
            if not row:
                raise HTTPException(404, "not found")
            return row
    finally:
        conn.close()


if __name__ == "__main__":
    import uvicorn
    uvicorn.run("order_producer:app", host="0.0.0.0", port=8080, workers=1, reload=False)
