#!/usr/bin/env python3
"""
Event Bus - 松耦合事件总线模块

用于 Agent 各组件之间的通信，实现发布/订阅模式。
"""

import asyncio
import logging
import time
from typing import Dict, Callable, List, Any, Optional
from dataclasses import dataclass, field
from enum import Enum

logger = logging.getLogger(__name__)


class EventType(str, Enum):
    """事件类型枚举"""
    # 对话事件
    DIALOGUE_START = "dialogue_start"
    DIALOGUE_END = "dialogue_end"
    USER_SPEECH = "user_speech"
    ASSISTANT_SPEECH = "assistant_speech"
    INTERRUPT = "interrupt"

    # 音频事件
    AUDIO_FRAME = "audio_frame"
    AUDIO_SEGMENT = "audio_segment"
    AUDIO_PLAYING = "audio_playing"
    AUDIO_STOPPED = "audio_stopped"

    # 记忆事件
    MEMORY_EXTRACTED = "memory_extracted"
    MEMORY_SAVED = "memory_saved"
    CONFLICT_DETECTED = "conflict_detected"

    # 视觉事件
    VISION_FRAME = "vision_frame"
    SCENE_IDENTIFIED = "scene_identified"

    # 系统事件
    SYSTEM_ERROR = "system_error"
    COMPONENT_READY = "component_ready"
    COMPONENT_ERROR = "component_error"
    HEALTH_CHECK = "health_check"


@dataclass
class EventBusEvent:
    """事件总线事件"""
    type: EventType
    data: Dict[str, Any] = field(default_factory=dict)
    timestamp: float = field(default_factory=time.time)
    source: str = "unknown"
    priority: int = 0  # 优先级，越高越先处理


class EventBus:
    """事件总线 - 支持同步和异步处理器"""
    
    def __init__(self, max_queue_size: int = 1000):
        self._handlers: Dict[EventType, List[Callable]] = {}
        self._async_handlers: Dict[EventType, List[Callable]] = {}
        self._queue: asyncio.Queue = asyncio.Queue(maxsize=max_queue_size)
        self._running = False
        self._task: Optional[asyncio.Task] = None
        self._stats = {
            "published": 0,
            "processed": 0,
            "errors": 0
        }
    
    def subscribe(self, event_type: EventType, handler: Callable, is_async: bool = False):
        """
        订阅事件
        
        Args:
            event_type: 事件类型
            handler: 处理函数（同步或异步）
            is_async: 是否异步处理
        """
        if is_async:
            if event_type not in self._async_handlers:
                self._async_handlers[event_type] = []
            self._async_handlers[event_type].append(handler)
            logger.debug(f"📡 异步订阅: {event_type} ← {handler.__name__}")
        else:
            if event_type not in self._handlers:
                self._handlers[event_type] = []
            self._handlers[event_type].append(handler)
            logger.debug(f"📡 同步订阅: {event_type} ← {handler.__name__}")
    
    def unsubscribe(self, event_type: EventType, handler: Callable):
        """取消订阅"""
        if event_type in self._handlers and handler in self._handlers[event_type]:
            self._handlers[event_type].remove(handler)
        if event_type in self._async_handlers and handler in self._async_handlers[event_type]:
            self._async_handlers[event_type].remove(handler)
    
    def publish(self, event_type: EventType, data: Dict[str, Any], 
                source: str = "unknown", priority: int = 0):
        """
        发布事件
        
        Args:
            event_type: 事件类型
            data: 事件数据
            source: 事件来源
            priority: 优先级
        """
        event = EventBusEvent(
            type=event_type,
            data=data,
            source=source,
            priority=priority
        )
        
        # 立即处理同步处理器
        if event_type in self._handlers:
            for handler in self._handlers[event_type]:
                try:
                    handler(event)
                except Exception as e:
                    logger.error(f"❌ 同步处理器错误: {handler.__name__} - {e}")
                    self._stats["errors"] += 1
        
        # 异步处理器入队
        if event_type in self._async_handlers:
            self._queue.put_nowait(event)
        
        self._stats["published"] += 1
        logger.debug(f"📨 发布事件: {event_type} (来源: {source})")
    
    async def start(self):
        """启动事件处理器"""
        self._running = True
        self._task = asyncio.create_task(self._process_queue())
        logger.info("✅ 事件总线已启动")
    
    async def stop(self):
        """停止事件处理器"""
        self._running = False
        if self._task:
            self._task.cancel()
            try:
                await self._task
            except asyncio.CancelledError:
                pass
        logger.info("⏹️ 事件总线已停止")
    
    async def _process_queue(self):
        """处理异步事件队列"""
        while self._running:
            try:
                event = await asyncio.wait_for(self._queue.get(), timeout=0.1)
                
                if event.type in self._async_handlers:
                    for handler in self._async_handlers[event.type]:
                        try:
                            await handler(event)
                            self._stats["processed"] += 1
                        except Exception as e:
                            logger.error(f"❌ 异步处理器错误: {handler.__name__} - {e}")
                            self._stats["errors"] += 1
                
                self._queue.task_done()
                
            except asyncio.TimeoutError:
                continue
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.error(f"❌ 事件队列处理错误: {e}")
                self._stats["errors"] += 1
    
    def get_stats(self) -> Dict[str, int]:
        """获取统计信息"""
        return self._stats.copy()
    
    def clear_handlers(self):
        """清空所有处理器"""
        self._handlers.clear()
        self._async_handlers.clear()
        logger.info("🧹 已清空所有事件处理器")


# ========== 便捷装饰器 ==========

def on_event(event_type: EventType, async_handler: bool = False):
    """
    事件处理器装饰器
    
    Usage:
        @on_event(EventType.USER_SPEECH, async_handler=True)
        async def handle_speech(event: EventBusEvent):
            pass
    """
    def decorator(func: Callable) -> Callable:
        func._event_handler = event_type
        func._async_handler = async_handler
        return func
    return decorator


# ========== 示例用法 ==========

if __name__ == "__main__":
    # 示例
    bus = EventBus()
    
    @on_event(EventType.USER_SPEECH, async_handler=True)
    async def handle_speech(event: EventBusEvent):
        print(f"🎤 收到语音: {event.data.get('text')}")
    
    @on_event(EventType.ASSISTANT_SPEECH, async_handler=True)
    async def handle_reply(event: EventBusEvent):
        print(f"🔊 回复: {event.data.get('text')}")
    
    async def main():
        await bus.start()
        
        # 发布事件
        bus.publish(EventType.USER_SPEECH, {"text": "你好"})
        bus.publish(EventType.ASSISTANT_SPEECH, {"text": "你好！有什么可以帮你的？"})
        
        await asyncio.sleep(1)
        await bus.stop()
        
        print(f"\n统计: {bus.get_stats()}")
    
    asyncio.run(main())
