#!/usr/bin/env python3
"""
Risk Consumer：风控消费者

业务规则（示例）：
- amount > 50000 → 标记可疑
- 同 user_id 1 分钟内 > 10 笔 → 标记可疑

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

import os
import time
from collections import defaultdict, deque
from confluent_kafka import Consumer, Producer
from confluent_kafka.serialization import SerializationContext, MessageField
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroDeserializer

BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
SR_URL = os.getenv("SR_URL", "http://localhost:8081")
TOPIC = "orders.events"
DLQ = "orders.events.dlq"
GROUP = "risk-group"

WINDOW_MS = 60_000
WINDOW_LIMIT = 10
AMOUNT_LIMIT = 50_000


def main():
    sr = SchemaRegistryClient({"url": SR_URL})
    deser = AvroDeserializer(sr, from_dict=lambda d, ctx: d)

    consumer = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
        "max.poll.interval.ms": 600_000,
        "max.poll.records": 100,
        "session.timeout.ms": 45000,
        "partition.assignment.strategy": "cooperative-sticky",
    })
    dlq_producer = Producer({"bootstrap.servers": BOOTSTRAP})
    consumer.subscribe([TOPIC])

    user_window = defaultdict(lambda: deque())
    print(f"[risk] start, group={GROUP}")

    try:
        while True:
            msg = consumer.poll(1.0)
            if msg is None:
                continue
            if msg.error():
                print(f"[risk] err: {msg.error()}")
                continue

            try:
                event = deser(msg.value(), SerializationContext(TOPIC, MessageField.VALUE))
            except Exception as e:
                print(f"[risk] decode fail → DLQ: {e}")
                dlq_producer.produce(DLQ, key=msg.key(), value=msg.value(),
                                     headers=[("error", str(e).encode())])
                consumer.commit(msg)
                continue

            # ===== 业务规则 =====
            risky = False
            reasons = []
            if event["amount"] > AMOUNT_LIMIT:
                risky = True
                reasons.append(f"amount {event['amount']} > {AMOUNT_LIMIT}")

            uid = event["user_id"]
            now = time.time() * 1000
            dq = user_window[uid]
            while dq and dq[0] < now - WINDOW_MS:
                dq.popleft()
            dq.append(now)
            if len(dq) > WINDOW_LIMIT:
                risky = True
                reasons.append(f"user {uid} {len(dq)} orders/min")

            tag = "🚨 RISKY" if risky else "✓ ok"
            print(f"[risk] {tag} order={event['order_id']} amt={event['amount']} reasons={reasons}")

            # ===== 提交 offset =====
            consumer.commit(msg, asynchronous=False)
    except KeyboardInterrupt:
        pass
    finally:
        consumer.close()


if __name__ == "__main__":
    main()
