# 当前流程 vs 松耦合架构对比

## 当前流程（紧密耦合）

```
┌──────────────────────────────────────────────────────────────┐
│                    run_worker.py / agent.py                   │
│                                                              │
│   LiveKit Job                                              │
│        │                                                     │
│        ▼                                                     │
│   VoiceAgent.initialize()                                   │
│        │                                                     │
│        ├── 加载 Config                                       │
│        ├── 初始化 STT (FunASR)                               │
│        ├── 初始化 LLM (Agnes AI)                             │
│        ├── 初始化 TTS (EdgeTTS)                              │
│        ├── 初始化 MemoryStore (ProfileManager) ← 直接依赖    │
│        └── 启动主循环                                         │
│                                                              │
│   主循环:                                                    │
│   ├─ STT 转录音频 → text                                     │
│   ├─ 直接调用 profile_mgr.add_fact()                         │
│   ├─ LLM.generate(text, context)                             │
│   └─ TTS.speak(response)                                     │
└──────────────────────────────────────────────────────────────┘
```

**问题：**
- Agent 直接依赖 ProfileManager
- 无法轻松替换组件
- 错误传播到主循环

---

## 松耦合架构（事件驱动）

```
┌──────────────────────────────────────────────────────────────┐
│                        EventBus                               │
│                     (消息总线)                                 │
│                                                              │
│   Events:                                                    │
│   • dialogue_event {text, speaker, timestamp}                │
│   • vision_event {image_bytes, timestamp}                    │
│   • user_action {action_type, params}                        │
│   • system_event {type, message}                             │
└──────────────────────────────────────────────────────────────┘
         ▲              ▲              ▲              ▲
         │              │              │              │
    ┌────┴────┐    ┌────┴────┐    ┌────┴────┐    ┌────┴────┐
    │  Agent  │    │ Memory  │    │  WeMM   │    │   Wiki  │
    │  Worker │    │  Store  │    │ Client  │    │  Service│
    └─────────┘    └─────────┘    └─────────┘    └─────────┘
    发布/订阅     订阅/处理      订阅/处理       订阅/处理
```

---

## 各组件职责与接口

### 1. Agent Worker（只负责语音对话）

```python
class AgentWorker:
    def __init__(self, config: Config, event_bus: EventBus):
        self.stt = FunASR(config.stt)
        self.llm = AgnesLLM(config.llm)
        self.tts = EdgeTTS(config.tts)
        self.event_bus = event_bus  # 只接收 EventBus
    
    async def run(self, room):
        # 1. STT 转录
        text = await self.stt.transcribe(audio_frame)
        
        # 2. 发布对话事件（不再直接调用 MemoryStore）
        self.event_bus.publish("dialogue_event", {
            "text": text,
            "speaker": "user",
            "timestamp": time.time()
        })
        
        # 3. LLM 生成回复
        response = await self.llm.generate(text)
        
        # 4. 发布回复事件
        self.event_bus.publish("dialogue_event", {
            "text": response,
            "speaker": "assistant",
            "timestamp": time.time()
        })
        
        # 5. TTS 合成
        await self.tts.speak(response)
```

---

### 2. Memory Store（独立处理记忆）

```python
class MemoryStoreHandler:
    def __init__(self, profile_mgr: ProfileManager, event_bus: EventBus):
        self.pm = profile_mgr
        self.event_bus = event_bus
        
        # 注册事件处理器
        event_bus.subscribe("dialogue_event", self.on_dialogue)
        event_bus.subscribe("conflict_detected", self.on_conflict)
    
    async def on_dialogue(self, event):
        """处理对话事件，自动提取事实"""
        text = event["data"]["text"]
        speaker = event["data"]["speaker"]
        
        # 提取事实
        facts = self.pm.extract_facts_from_dialogue(text, speaker)
        
        # 写入记忆
        for fact in facts:
            self.pm.add_fact(speaker, fact["predicate"], fact["object"])
        
        # 检测矛盾
        conflicts = self.pm.detect_conflicts(speaker)
        if conflicts:
            self.event_bus.publish("conflict_warning", {
                "speaker": speaker,
                "conflicts": conflicts
            })
```

---

### 3. WeMM Client（向量生成服务）

```python
class WeMMEmbeddingService:
    def __init__(self, config: dict, event_bus: EventBus):
        self.api_url = config.get("wemm_api_url")
        self.dimension = config.get("dimension", 768)
        self.event_bus = event_bus
        
        # 注册事件处理器
        event_bus.subscribe("vision_event", self.on_vision)
    
    async def on_vision(self, event):
        """处理视觉事件，生成向量"""
        image_bytes = event["data"]["image_bytes"]
        text_desc = event["data"].get("text_description", "")
        
        # 生成向量
        if text_desc:
            vector = await self.embed_text(text_desc)
        else:
            vector = await self.embed_image(image_bytes)
        
        # 存储向量
        await self.store_embedding(vector, event["data"])
    
    async def embed_text(self, text: str) -> List[float]:
        # 调用 WeMM API
        return await self.client.encode_text(text)
    
    async def embed_image(self, image_bytes: bytes) -> List[float]:
        return await self.client.encode_image(image_bytes)
```

---

### 4. Knowledge Wiki（RAG 检索）

```python
class KnowledgeWikiService:
    def __init__(self, config: dict, event_bus: EventBus, vector_service: WeMMEmbeddingService):
        self.config = config
        self.vector_service = vector_service
        self.event_bus = event_bus
        
        # 注册事件处理器
        event_bus.subscribe("query_event", self.on_query)
    
    async def on_query(self, event):
        """处理查询事件，返回相关知识"""
        query = event["data"]["query"]
        
        # 向量化查询
        query_vec = await self.vector_service.embed_text(query)
        
        # 相似度搜索
        results = await self.wiki_store.search(query_vec, top_k=3)
        
        # 返回结果
        return results
```

---

## 启动流程（解耦版）

```python
# main.py - 组装各组件
async def main():
    # 1. 加载配置
    config = Config.load("config.yaml")
    
    # 2. 创建事件总线
    event_bus = EventBus()
    
    # 3. 初始化各组件（传入 event_bus）
    memory_handler = MemoryStoreHandler(
        profile_mgr=ProfileManager(),
        event_bus=event_bus
    )
    
    wemm_service = WeMMEmbeddingService(
        config=config.vision,
        event_bus=event_bus
    )
    
    wiki_service = KnowledgeWikiService(
        config=config,
        event_bus=event_bus,
        vector_service=wemm_service
    )
    
    # 4. 初始化 Agent（传入 event_bus）
    agent = AgentWorker(
        config=config,
        event_bus=event_bus
    )
    
    # 5. 启动各服务
    await agent.start()
    await memory_handler.start()
    await wemm_service.start()
    await wiki_service.start()
    
    logger.info("✅ 所有组件已启动，进入事件驱动模式")
```

---

## 数据流示例

```
用户: "我叫张三，我喜欢蓝色"
         │
         ▼
┌─────────────────┐
│   Agent Worker  │  STT 转录
└────────┬────────┘
         │
         │ publish("dialogue_event", {text: "我叫张三...", speaker: "user"})
         ▼
┌─────────────────────────────────────────────────┐
│                 EventBus                        │
└────────┬──────────────────────────┬──────────────┘
         │                          │
         ▼                          ▼
┌─────────────────┐      ┌─────────────────┐
│ Memory Store    │      │  Agent Worker   │
│ Handler         │      │                 │
│                 │      │  LLM 生成回复   │
│ 1. 提取事实     │      │                 │
│ 2. 写入记忆     │      │ publish("reply")│
│ 3. 检测矛盾     │      └────────┬────────┘
└─────────────────┘               │
                                  ▼
                           ┌─────────────────┐
                           │    TTS Output   │
                           └─────────────────┘
```

---

## 替换组件示例

### 场景1：更换 STT 引擎

```python
# 只需替换 AgentWorker 的 stt 实现
agent = AgentWorker(
    config=config,
    event_bus=event_bus,
    stt=WhisperSTT(config.stt)  # 换成 Whisper
)
```

### 场景2：更换记忆存储

```python
# 只需替换 MemoryStoreHandler 的实现
memory_handler = MySQLMemoryHandler(
    db_config=config.mysql,
    event_bus=event_bus
)
```

### 场景3：关闭记忆功能

```python
# 不注册 MemoryStoreHandler 即可
# event_bus 仍然工作，只是没有记忆组件
```

---

## 配置对比

| 项目 | 当前（紧密耦合） | 新架构（松耦合） |
|------|----------------|----------------|
| 组件依赖 | Agent 直接 import MemoryStore | 通过 EventBus 解耦 |
| 错误传播 | 单个组件失败影响主流程 | 各组件独立 try-catch |
| 替换组件 | 需改代码 | 只需改配置 |
| 测试难度 | 需初始化所有依赖 | 可单独测试每个组件 |
| 扩展性 | 难 | 易（新增事件类型即可） |
