refactor: 重构目录结构 - 简化层级
Some checks failed
构建并部署 AI Agent 服务 / deploy (push) Has been cancelled

This commit is contained in:
2026-04-29 12:52:41 +08:00
parent 223d1c9afd
commit ef5113bffb
54 changed files with 42 additions and 1819 deletions

View File

@@ -0,0 +1 @@
"""子图模块"""

View File

@@ -0,0 +1,391 @@
# 通讯录子图 (Contact Subgraph)
该子图负责处理通讯录管理、邮件读取与发送等功能,基于 LangGraph 状态机编排多阶段工作流,支持联系人 CRUD、IMAP 邮箱绑定、邮件审核发送等核心能力。子图设计遵循"安全优先、审核强制、隐私保护"原则,通过敏感信息加密和人工审核保障数据安全。
> **使用公共工具**意图理解、人工审核、格式化输出、检查点持久化、条件路由、LLM 调用、数据库工具、状态基类
---
## 🎯 核心架构
### 技术栈
| 层级 | 组件 | 说明 |
|:-----|:-----|:-----|
| **编排框架** | LangGraph StateGraph | 状态机驱动的子图工作流编排,支持中断恢复 |
| **LLM 服务** | 智谱 AI / DeepSeek API | 意图理解、邮件草稿生成、联系人信息提取(使用公共 LLM 工具) |
| **关系存储** | PostgreSQL | 联系人、邮箱配置持久化(使用公共数据库工具) |
| **邮件协议** | imaplib / smtplib | 邮件读取与发送 |
| **加密存储** | cryptography | 邮箱密码等敏感信息加密 |
### 子图分层架构
```
┌─────────────────────────────────────────────────────────────────┐
│ 主图 (Main Graph) │
└──────────────────────────────┬──────────────────────────────────┘
│ 状态映射 / 结果聚合
┌─────────────────────────────────────────────────────────────────┐
│ 通讯录子图接口层 │
│ - 状态转换:主状态 ↔ 子图状态(使用公共状态基类) │
│ - 错误传播与优雅降级 │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 工作流编排层 │
│ - 节点调度与条件路由(使用公共路由工具) │
│ - 人工审核节点暂停/恢复管理(使用公共审核工具) │
│ - 状态持久化与检查点(使用公共检查点工具) │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 节点层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │
│ │意图理解 │ │联系人CRUD│ │邮件读取 │ │草稿生成 │ │人工审核│ │
│ │(公共工具)│ │ │ │ │ │ │ │(公共) │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │邮件发送 │ │智能嗅探 │ │格式输出 │ │
│ │ │ │ │ │(公共) │ │
│ └──────────┘ └──────────┘ └──────────┘ │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 工具层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │数据库工具│ │IMAP工具 │ │SMTP工具 │ │加密工具 │ │
│ │(公共) │ │ │ │ │ │ │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────────┘
```
### 数据流总览
通讯录子图根据意图类型分支执行,关键操作前强制人工审核。
```
用户请求
┌─────────────┐
│ 意图理解 │ ← 使用公共意图理解工具
└──────┬──────┘
├──────────┬──────────┬──────────┬──────────┐
▼ ▼ ▼ ▼ ▼
联系人CRUD 邮件读取 邮件发送 智能嗅探 列表查询
│ │ │ │ │
▼ ▼ ▼ ▼ ▼
数据库操作 IMAP读取 草稿生成 信息提取 结果展示
│ │ │ │ │
│ │ ▼ │ │
│ │ 人工审核 │ │
│ │ (公共) │ │
│ │ │ │ │
│ │ ▼ │ │
│ │ SMTP发送 │ │
│ │ │ │ │
└──────────┴──────────┴──────────┴──────────┘
格式输出 ← 使用公共格式化工具
返回主图
```
---
## 📂 模块与文件结构
```
app/contact/
├── __init__.py
├── graph.py # 子图构建入口,定义状态图与路由
├── state.py # 子图状态定义(继承公共状态基类)
├── nodes/ # 节点实现
│ ├── __init__.py
│ ├── crud.py # 联系人CRUD节点
│ ├── email_read.py # 邮件读取节点
│ ├── draft.py # 邮件草稿生成节点
│ ├── email_send.py # 邮件发送节点
│ └── sniff.py # 智能嗅探节点
├── tools/ # 子图特有工具集
│ ├── imap.py # IMAP邮件读取工具
│ ├── smtp.py # SMTP邮件发送工具
│ └── crypto.py # 敏感信息加密工具
└── persistence/ # (使用公共检查点工具,无需单独实现)
```
> **注意**:以下模块使用公共工具,无需单独实现:
> - 意图理解节点 → 使用 `agent_subgraphs.common.intent`
> - 人工审核节点 → 使用 `agent_subgraphs.common.human_loop`
> - 格式输出节点 → 使用 `agent_subgraphs.common.format`
> - 检查点持久化 → 使用 `agent_subgraphs.common.checkpoint`
> - 条件路由 → 使用 `agent_subgraphs.common.routing`
> - LLM 调用 → 使用 `agent_subgraphs.common.llm`
> - 数据库操作 → 使用 `agent_subgraphs.common.db`
---
## 🎯 演进路线与核心机制
### Level 1基础联系人管理
**核心机制**:联系人 CRUD 操作 + 基础列表查询。
- 使用 LLM 解析用户请求,提取联系人信息(姓名、电话、邮箱等)。
- 数据库存储联系人信息,支持增删改查。
- 支持按姓名模糊查询和列表展示。
**适用场景**:保存联系人、查询电话、修改信息等基础操作。
**实现指引**:意图理解节点识别 CRUD 类型,路由到对应节点。
### Level 2邮件读取与基础发送
**核心机制**IMAP 邮箱绑定 + 邮件列表读取 + 简单发送。
- 支持配置多个 IMAP/SMTP 邮箱账户。
- 读取收件箱邮件列表,展示主题、发件人、时间。
- 生成邮件草稿,人工审核后发送。
**适用场景**:查收邮件、发送简单邮件。
**实现指引**:邮箱密码加密存储,发送前强制人工审核。
### Level 3智能嗅探与上下文感知
**核心机制**:对话中自动识别联系人信息,主动询问是否保存。
- 在任意对话中监听"人名+联系方式"模式。
- 识别到潜在联系人信息时,主动询问用户是否保存。
- 支持确认后自动创建联系人记录。
**适用场景**:日常对话中自然保存联系人。
**实现指引**:与主图协作,在对话节点后插入嗅探检查点。
### Level 4邮件智能处理
**核心机制**:邮件内容理解 + 智能回复建议 + 批量处理。
- 读取邮件全文,理解内容意图。
- 基于邮件内容生成智能回复草稿。
- 支持邮件分类、标记、归档操作。
**适用场景**:邮件管理、智能回复。
**实现指引**:利用 LLM 进行邮件内容理解和回复生成。
### Level 5通讯录智能助理
**核心机制**:联系人关系图谱 + 智能推荐 + 多模态交互。
- 构建联系人关系网络,识别社交圈子。
- 基于历史交互推荐联系人。
- 支持名片扫描、语音输入等多模态交互。
**适用场景**:深度联系人管理、智能社交助理。
---
## 🔧 核心组件详解
### 1. 意图理解节点
**职责**:接收用户原始请求,区分联系人 CRUD、邮件读取、邮件发送、智能嗅探等意图类型。
**输入**:用户自然语言请求。
**输出**
- `intent_type`意图类别枚举contact_crud / email_read / email_send / sniff / list
- `extracted_info`:提取的关键信息(联系人姓名、邮件主题等)。
**实现要点**
- 使用 LLM 进行少样本分类,输出结构化 JSON。
- 关键词匹配兜底(如"保存"、"添加" → contact_crud
### 2. 联系人 CRUD 节点
**职责**:执行联系人的增删改查操作,与数据库交互。
**输入**:意图类型、提取的联系人信息。
**输出**
- `operation_result`:操作结果(成功/失败)。
- `contact_data`:操作后的联系人数据。
**实现要点**
- 使用 SQLAlchemy 与 PostgreSQL 交互。
- 支持部分字段更新(如只修改电话)。
- 删除操作前二次确认。
### 3. 邮件读取节点
**职责**:通过 IMAP 协议连接邮箱,读取邮件列表或单封邮件详情。
**输入**:邮箱配置、读取指令(列表/详情)。
**输出**
- `email_list`:邮件列表(主题、发件人、时间)。
- `email_content`:单封邮件详情(如需要)。
**实现要点**
- 支持配置多个邮箱账户。
- 密码使用 cryptography 加密存储。
- 分页读取邮件列表,避免一次性加载过多。
### 4. 邮件草稿生成节点
**职责**:根据用户指令生成邮件草稿,包含收件人、主题、正文。
**输入**:用户发送邮件的指令、上下文。
**输出**
- `draft_recipient`:收件人邮箱。
- `draft_subject`:邮件主题。
- `draft_body`:邮件正文。
**实现要点**
- 使用 LLM 生成自然流畅的邮件内容。
- 从通讯录智能匹配收件人邮箱。
- 支持多语言邮件生成。
### 5. 人工审核节点
**职责**:在邮件发送、联系人删除等关键操作前挂起,等待用户确认。
**输入**:待审核内容(邮件草稿、删除确认等)。
**输出**
- `user_approved`:是否通过审核。
- `user_modification`:用户修改内容(如有)。
**实现要点**
- 使用 LangGraph `interrupt` 机制实现挂起。
- 支持用户修改草稿后再次审核。
- 超时自动取消操作。
### 6. 邮件发送节点
**职责**:审核通过后,通过 SMTP 协议发送邮件。
**输入**:审核通过的邮件草稿。
**输出**
- `send_result`:发送结果。
- `sent_time`:发送时间。
**实现要点**
- 使用 smtplib 发送邮件。
- 支持抄送、密送。
- 发送失败重试机制。
### 7. 智能嗅探节点
**职责**:在对话中检测联系人信息,主动询问是否保存。
**输入**:用户对话历史。
**输出**
- `detected_contact`:检测到的潜在联系人信息。
- `should_ask`:是否应该询问用户。
**实现要点**
- 使用 NER 识别实体(人名、电话、邮箱)。
- 与已有联系人去重。
- 询问用户确认后自动保存。
---
## 🔀 条件路由详解
### 入口路由:意图分支
- **位置**:意图理解节点之后。
- **条件**
- `intent_type == "contact_crud"` → 联系人 CRUD 节点。
- `intent_type == "email_read"` → 邮件读取节点。
- `intent_type == "email_send"` → 草稿生成节点。
- `intent_type == "sniff"` → 智能嗅探节点。
- `intent_type == "list"` → 列表查询节点。
### 审核路由
- **位置**:人工审核节点之后。
- **条件**
- `user_approved == True` → 执行操作(发送/删除等)。
- `user_approved == False``user_modification` 存在 → 返回草稿生成节点。
- `user_approved == False` 且无修改 → 取消操作。
### 嗅探路由
- **位置**:智能嗅探节点之后。
- **条件**
- `should_ask == True` → 询问用户确认。
- `should_ask == False` → 直接结束,不打扰用户。
---
## 📊 状态设计
### 状态结构概览
| 分组 | 字段 | 类型 | 说明 |
|:-----|:-----|:-----|:-----|
| **输入** | `user_input` | `str` | 用户原始请求 |
| **意图** | `intent_type` | `str` | 意图类别 |
| | `extracted_info` | `dict` | 提取的关键信息 |
| **联系人** | `contact_data` | `dict` | 联系人数据 |
| | `contact_list` | `list[dict]` | 联系人列表 |
| **邮件** | `email_config` | `dict` | 邮箱配置 |
| | `email_list` | `list[dict]` | 邮件列表 |
| | `draft_recipient` | `str` | 草稿收件人 |
| | `draft_subject` | `str` | 草稿主题 |
| | `draft_body` | `str` | 草稿正文 |
| **审核** | `pending_action` | `dict` | 待审核操作 |
| | `user_approved` | `bool` | 是否通过审核 |
| | `user_modification` | `dict` | 用户修改 |
| **控制流** | `current_phase` | `str` | 当前执行阶段 |
| | `next_node` | `str` | 下一节点名称 |
| | `interrupt_point` | `str` | 中断点标识 |
| **输出** | `final_result` | `str` | 最终结果 |
---
## 🔄 工作流程与中断恢复
### 邮件发送流程
| 步骤 | 节点 | 人工干预 |
|:-----|:-----|:---------|
| 1 | 意图理解 | 否 |
| 2 | 草稿生成 | 否 |
| 3 | 人工审核 | **是** |
| 4 | 邮件发送 | 否 |
| 5 | 格式输出 | 否 |
### 联系人保存流程
| 步骤 | 节点 | 人工干预 |
|:-----|:-----|:---------|
| 1 | 意图理解 | 否 |
| 2 | 联系人 CRUD | 否 |
| 3 | 格式输出 | 否 |
### 智能嗅探流程
| 步骤 | 节点 | 人工干预 |
|:-----|:-----|:---------|
| 1 | 智能嗅探 | 否 |
| 2 | 询问确认 | **是** |
| 3 | 联系人 CRUD | 否 |
### 中断恢复机制
子图支持在人工审核节点中断,恢复时从审核点继续执行。

View File

@@ -0,0 +1,52 @@
"""
通讯录子图 - 完善版
Contact Subgraph Module - Complete
"""
from .state import (
ContactState,
Contact,
Email,
ContactAction
)
from .graph import build_contact_subgraph
from .nodes import (
parse_intent,
list_contacts,
add_contact,
list_emails,
generate_email_draft,
human_review,
send_email,
sniff_contacts,
format_result,
should_continue
)
from .api_client import contact_api, ContactAPIClient
__all__ = [
# State
"ContactState",
"Contact",
"Email",
"ContactAction",
# Graph
"build_contact_subgraph",
# Nodes
"parse_intent",
"list_contacts",
"add_contact",
"list_emails",
"generate_email_draft",
"human_review",
"send_email",
"sniff_contacts",
"format_result",
"should_continue",
# API
"contact_api",
"ContactAPIClient"
]

View File

@@ -0,0 +1,286 @@
"""
通讯录子图 API 调用工具
支持模拟数据和真实数据库两种模式
"""
from typing import Dict, Any, Optional, List
from datetime import datetime
from dataclasses import dataclass
from .state import Contact, Email
# ========== 模拟数据(保留作为备选)==========
# 模拟数据库
MOCK_CONTACTS_DB = {}
MOCK_EMAILS_DB = []
@dataclass
class ContactAPIClient:
"""
通讯录 API 客户端 - 支持真实数据库和模拟模式
使用方式:
1. 真实数据库模式:传入 conn 参数
2. 模拟模式:不传入 conn或 conn 为 None
"""
def __init__(self, conn=None):
"""
初始化
Args:
conn: 数据库连接(来自 checkpointer.conn为 None 时使用模拟模式
"""
self.conn = conn
self._use_db = conn is not None
if self._use_db:
try:
from ...db.models import ContactRepository, ContactEntity
self._repo = ContactRepository(conn)
except Exception as e:
print(f"Repository 初始化失败,回退到模拟模式: {e}")
self._use_db = False
self._repo = None
# ========== 真实数据库方法 ==========
async def list_contacts_db(self, user_id: str = "default") -> List[Contact]:
"""真实数据库:获取联系人列表"""
if not self._repo:
return await self.list_contacts_mock(user_id)
entities = await self._repo.list_by_user(user_id)
return [
Contact(
id=e.id,
name=e.name,
phone=e.phone,
email=e.email,
company=e.company,
position=e.position,
created_at=e.created_at
)
for e in entities
]
async def add_contact_db(self, user_id: str, contact: Contact) -> bool:
"""真实数据库:添加联系人"""
if not self._repo:
return await self.save_contact_mock(user_id, contact)
from ...db.models import ContactEntity
entity = ContactEntity(
user_id=user_id,
name=contact.name,
phone=contact.phone,
email=contact.email,
company=contact.company,
position=contact.position,
created_at=contact.created_at or datetime.now().isoformat()
)
await self._repo.insert(entity)
return True
# ========== 模拟数据方法(保留)==========
def list_contacts_mock(self, user_id: str = "default") -> List[Contact]:
"""模拟查询联系人列表"""
if user_id not in MOCK_CONTACTS_DB:
# 初始化一些示例数据
MOCK_CONTACTS_DB[user_id] = [
Contact(
id="1",
name="张三",
phone="13800138000",
email="zhangsan@example.com",
company="科技公司",
position="工程师",
created_at=datetime.now().isoformat()
),
Contact(
id="2",
name="李四",
phone="13900139000",
email="lisi@example.com",
company="贸易公司",
position="经理",
created_at=datetime.now().isoformat()
),
Contact(
id="3",
name="王五",
phone="13700137000",
email="wangwu@example.com",
company="咨询公司",
position="顾问",
created_at=datetime.now().isoformat()
),
]
return MOCK_CONTACTS_DB[user_id]
def extract_contact_info_mock(self, query: str) -> Optional[Dict[str, Any]]:
"""模拟从查询中提取联系人信息"""
import re
# 提取邮箱
email_match = re.search(r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}', query)
# 提取手机号
phone_match = re.search(r'1[3-9]\d{9}', query)
# 提取姓名(简单匹配)
if any(keyword in query for keyword in ["添加", "add"]):
name = "未知"
clean_query = query
if email_match:
clean_query = clean_query.replace(email_match.group(), "")
if phone_match:
clean_query = clean_query.replace(phone_match.group(), "")
clean_query = clean_query.replace("添加", "").replace("add", "").replace("联系人", "").strip()
if clean_query:
name = clean_query
return {
"name": name,
"phone": phone_match.group() if phone_match else "",
"email": email_match.group() if email_match else "",
"created_at": datetime.now().isoformat()
}
return None
def save_contact_mock(self, user_id: str, contact: Contact) -> bool:
"""模拟保存联系人"""
if user_id not in MOCK_CONTACTS_DB:
MOCK_CONTACTS_DB[user_id] = []
if not contact.id:
contact.id = str(len(MOCK_CONTACTS_DB[user_id]) + 1)
MOCK_CONTACTS_DB[user_id].append(contact)
return True
def list_emails_mock(self) -> List[Email]:
"""模拟查询邮件列表"""
global MOCK_EMAILS_DB
if not MOCK_EMAILS_DB:
MOCK_EMAILS_DB = [
Email(
id="1",
subject="会议邀请AI 技术分享",
sender="admin@example.com",
recipients=["user@example.com"],
date=datetime.now().isoformat(),
body="你好,下周一将举办 AI 技术分享会,欢迎参加。"
),
Email(
id="2",
subject="项目进度更新",
sender="manager@example.com",
recipients=["user@example.com"],
date=datetime.now().isoformat(),
body="项目进度良好,继续保持。"
),
]
return MOCK_EMAILS_DB
def generate_email_draft_mock(self, query: str) -> Dict[str, str]:
"""模拟生成邮件草稿"""
return {
"subject": f"Re: {query}",
"recipient": "recipient@example.com",
"body": "你好,\n\n这是一封自动生成的邮件草稿。\n\n此致,\n你的助手"
}
def send_email_mock(self, recipient: str, subject: str, body: str) -> Dict[str, Any]:
"""模拟发送邮件"""
global MOCK_EMAILS_DB
MOCK_EMAILS_DB.append(
Email(
id=str(len(MOCK_EMAILS_DB) + 1),
subject=subject,
sender="me@example.com",
recipients=[recipient],
date=datetime.now().isoformat(),
body=body
)
)
return {
"success": True,
"message": "邮件发送成功"
}
def sniff_contacts_mock(self, query: str) -> Dict[str, Any]:
"""模拟智能嗅探联系人"""
import re
emails = re.findall(r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}', query)
phones = re.findall(r'1[3-9]\d{9}', query)
contacts = []
for i, email in enumerate(emails):
contacts.append({
"name": f"联系人{i+1}",
"email": email,
"phone": phones[i] if i < len(phones) else ""
})
return {
"contacts": contacts,
"count": len(contacts),
"suggestion": "是否添加这些联系人?"
}
# ========== 公共方法(自动选择模式)==========
async def list_contacts(self, user_id: str = "default") -> List[Contact]:
"""获取联系人列表(自动选择数据库或模拟模式)"""
if self._use_db:
return await self.list_contacts_db(user_id)
return self.list_contacts_mock(user_id)
async def add_contact(self, user_id: str, contact: Contact) -> bool:
"""添加联系人(自动选择数据库或模拟模式)"""
if self._use_db:
return await self.add_contact_db(user_id, contact)
return self.save_contact_mock(user_id, contact)
async def list_emails(self, user_id: str = "default") -> List[Email]:
"""查询邮件列表(目前用模拟)"""
return self.list_emails_mock()
async def generate_email_draft(self, query: str) -> Dict[str, str]:
"""生成邮件草稿(目前用模拟)"""
return self.generate_email_draft_mock(query)
async def send_email(self, user_id: str, recipient: str, subject: str, body: str) -> bool:
"""发送邮件(目前用模拟)"""
result = self.send_email_mock(recipient, subject, body)
return result.get("success", False)
async def sniff_contacts(self, query: str) -> List[Contact]:
"""智能嗅探联系人(目前用模拟)"""
result = self.sniff_contacts_mock(query)
contact_dicts = result.get("contacts", [])
return [
Contact(
id=str(i+1),
name=c.get("name", ""),
phone=c.get("phone", ""),
email=c.get("email", ""),
company="",
position="",
created_at=datetime.now().isoformat()
)
for i, c in enumerate(contact_dicts)
]
# 全局实例(模拟模式,保留向后兼容)
contact_api = ContactAPIClient()

View File

@@ -0,0 +1,109 @@
"""
通讯录子图构建器
Contact Subgraph Builder
支持 API 注入的工厂模式
"""
from app.main_graph.graph import StateGraph, START, END
from .state import ContactState
from .nodes import create_contact_nodes
def build_contact_subgraph(contact_api=None):
"""
构建通讯录子图(工厂模式)
Args:
contact_api: 可选的 ContactAPIClient 实例(支持真实数据库或模拟模式)
不传入则使用默认模拟 API向后兼容
Returns:
配置好的 StateGraph
"""
# 创建节点(传入 API
nodes = create_contact_nodes(contact_api) if contact_api else None
# 如果没有传入 API使用向后兼容的导入
if nodes is None:
from .nodes import (
parse_intent,
list_contacts,
add_contact,
list_emails,
generate_email_draft,
human_review,
send_email,
sniff_contacts,
format_result,
should_continue
)
else:
parse_intent = nodes["parse_intent"]
list_contacts = nodes["list_contacts"]
add_contact = nodes["add_contact"]
list_emails = nodes["list_emails"]
generate_email_draft = nodes["generate_email_draft"]
human_review = nodes["human_review"]
send_email = nodes["send_email"]
sniff_contacts = nodes["sniff_contacts"]
format_result = nodes["format_result"]
should_continue = nodes["should_continue"]
# 创建图
graph = StateGraph(ContactState)
# 添加节点
graph.add_node("parse_intent", parse_intent)
graph.add_node("list_contacts", list_contacts)
graph.add_node("add_contact", add_contact)
graph.add_node("list_emails", list_emails)
graph.add_node("generate_email_draft", generate_email_draft)
graph.add_node("human_review", human_review)
graph.add_node("send_email", send_email)
graph.add_node("sniff_contacts", sniff_contacts)
graph.add_node("format_result", format_result)
# 添加边
# 从START开始
graph.add_edge(START, "parse_intent")
# 从parse_intent根据条件路由
graph.add_conditional_edges(
"parse_intent",
should_continue,
{
"list_contacts": "list_contacts",
"add_contact": "add_contact",
"list_emails": "list_emails",
"generate_email_draft": "generate_email_draft",
"sniff_contacts": "sniff_contacts",
}
)
# 从各个操作节点到format_result
graph.add_edge("list_contacts", "format_result")
graph.add_edge("add_contact", "format_result")
graph.add_edge("list_emails", "format_result")
graph.add_edge("sniff_contacts", "format_result")
# 邮件发送的特殊流程
graph.add_edge("generate_email_draft", "human_review")
# 从human_review根据条件路由
graph.add_conditional_edges(
"human_review",
should_continue,
{
"send_email": "send_email",
"format_result": "format_result",
}
)
# 发送邮件后到格式化
graph.add_edge("send_email", "format_result")
# 最终到END
graph.add_edge("format_result", END)
return graph

View File

@@ -0,0 +1,278 @@
"""
通讯录子图节点 - 使用公共工具版本
Contact Subgraph Nodes - Using Common Tools
支持 async 和 API 注入
"""
from typing import Dict, Any
from datetime import datetime
# 公共工具
from ..common import MarkdownFormatter
from .state import ContactState, ContactAction, Contact, Email
from .api_client import ContactAPIClient
# 模拟联系人数据库(临时存储,保留作为备选)
CONTACT_DB = {}
def create_contact_nodes(contact_api: ContactAPIClient):
"""
创建通讯录子图节点工厂函数
Args:
contact_api: 已初始化的 ContactAPIClient支持真实数据库或模拟模式
Returns:
节点函数字典
"""
async def parse_intent(state: ContactState) -> ContactState:
"""
解析用户意图节点
"""
query_lower = state.user_query.lower()
if any(keyword in query_lower for keyword in ["添加", "add", "新建", "save"]):
state.action = ContactAction.CONTACT_ADD
elif any(keyword in query_lower for keyword in ["联系人", "contact", "list"]):
state.action = ContactAction.CONTACT_LIST
state.action_params = {"query": state.user_query}
elif any(keyword in query_lower for keyword in ["邮件", "email", "inbox"]):
state.action = ContactAction.EMAIL_LIST
elif any(keyword in query_lower for keyword in ["发送邮件", "send email", "发邮件"]):
state.action = ContactAction.EMAIL_SEND
else:
state.action = ContactAction.SNIFF_CONTACTS
return state
async def list_contacts(state: ContactState) -> ContactState:
"""
列出联系人节点
"""
state.current_phase = "executing"
# 使用 API 客户端async
contacts = await contact_api.list_contacts(state.user_id)
state.contacts = contacts
return state
async def add_contact(state: ContactState) -> ContactState:
"""
添加联系人节点
"""
state.current_phase = "executing"
# 使用 API 客户端(简化添加,实际项目应解析用户输入)
new_contact = Contact(
id=str(len(CONTACT_DB) + 1),
name="新联系人",
email="new@example.com",
phone="13800000000",
created_at=datetime.now().isoformat()
)
# 保存到数据库
await contact_api.add_contact(state.user_id, new_contact)
state.current_contact = new_contact
return state
async def list_emails(state: ContactState) -> ContactState:
"""
列出邮件节点
"""
state.current_phase = "executing"
# 使用 API 客户端async
emails = await contact_api.list_emails(state.user_id)
state.emails = emails
return state
async def generate_email_draft(state: ContactState) -> ContactState:
"""
生成邮件草稿节点
"""
state.current_phase = "executing"
# 使用 API 客户端async
draft = await contact_api.generate_email_draft(state.user_query)
state.draft_recipient = draft.get("recipient", "recipient@example.com")
state.draft_subject = draft.get("subject", "邮件主题")
state.draft_body = draft.get("body", "邮件正文")
return state
async def sniff_contacts(state: ContactState) -> ContactState:
"""
嗅探联系人节点
"""
state.current_phase = "executing"
# 使用 API 客户端async
contacts = await contact_api.sniff_contacts(state.user_query)
state.sniffed_contacts = contacts
return state
async def format_result(state: ContactState) -> ContactState:
"""
格式化结果节点(使用公共工具)
"""
state.current_phase = "formatting"
md = MarkdownFormatter()
output_lines = []
output_lines.append("┌───────────────────────────────────┐")
output_lines.append("│ 📇 通讯录助手 │")
output_lines.append("└───────────────────────────────────┘")
output_lines.append("")
if state.action == ContactAction.CONTACT_LIST and state.contacts:
output_lines.append(md.heading("📇 联系人列表", 2))
output_lines.append("")
contact_data = [
{"姓名": c.name, "邮箱": c.email, "电话": c.phone or "-"}
for c in state.contacts
]
output_lines.append(md.table(contact_data))
elif state.action == ContactAction.EMAIL_LIST and state.emails:
output_lines.append(md.heading("📬 邮件列表", 2))
output_lines.append("")
# 兼容两种 date 格式
email_data = []
for e in state.emails:
date_str = e.date
if hasattr(e, 'received_at') and e.received_at:
try:
date_str = e.received_at.strftime('%Y-%m-%d %H:%M')
except:
pass
email_data.append({
"发件人": e.sender,
"主题": e.subject,
"时间": date_str
})
output_lines.append(md.table(email_data))
elif state.action == ContactAction.EMAIL_SEND and state.draft_subject:
output_lines.append(md.heading("📝 邮件草稿", 2))
output_lines.append("")
output_lines.append(f"**收件人**: {state.draft_recipient}")
output_lines.append(f"**主题**: {state.draft_subject}")
output_lines.append("")
output_lines.append(md.quote(state.draft_body))
elif state.action == ContactAction.SNIFF_CONTACTS and state.sniffed_contacts:
output_lines.append(md.heading("🔍 发现的联系人", 2))
output_lines.append("")
contact_data = [
{"姓名": c.name, "邮箱": c.email}
for c in state.sniffed_contacts
]
output_lines.append(md.table(contact_data))
else:
output_lines.append(md.heading("✨ 操作完成", 2))
output_lines.append("您的请求已处理。")
# 页脚提示
output_lines.append("")
output_lines.append("---")
output_lines.append("💡 提示:您可以继续查询联系人、查看邮件,或者生成邮件草稿!")
state.final_result = "\n".join(output_lines)
state.success = True
state.current_phase = "completed"
return state
async def human_review(state: ContactState) -> ContactState:
"""
人工审核节点(用于邮件草稿)
"""
state.current_phase = "reviewing"
# 标记需要审核,等待用户决定
state.needs_approval = True
return state
async def send_email(state: ContactState) -> ContactState:
"""
发送邮件节点
"""
state.current_phase = "executing"
# 使用 API 客户端发送邮件async
success = await contact_api.send_email(
state.user_id,
state.draft_recipient,
state.draft_subject,
state.draft_body
)
state.success = success
return state
def should_continue(state: ContactState) -> str:
"""
条件路由函数:根据 action 和状态决定下一个节点
"""
# 如果是从 human_review 来的,根据审核状态决定
if state.current_phase == "reviewing":
if state.needs_approval:
# 这里会等待用户操作,实际运行时通过 checkpointer 或后端 API 处理
return "format_result"
else:
return "send_email"
# 普通路由
action = state.action
if action == ContactAction.CONTACT_LIST:
return "list_contacts"
elif action == ContactAction.CONTACT_ADD:
return "add_contact"
elif action == ContactAction.EMAIL_LIST:
return "list_emails"
elif action == ContactAction.EMAIL_SEND:
return "generate_email_draft"
elif action == ContactAction.SNIFF_CONTACTS:
return "sniff_contacts"
else:
return "format_result"
# 返回节点字典
return {
"parse_intent": parse_intent,
"list_contacts": list_contacts,
"add_contact": add_contact,
"list_emails": list_emails,
"generate_email_draft": generate_email_draft,
"sniff_contacts": sniff_contacts,
"format_result": format_result,
"human_review": human_review,
"send_email": send_email,
"should_continue": should_continue
}
# ========== 向后兼容的全局版本(使用模拟 API ==========
from .api_client import contact_api as _default_contact_api
# 创建默认节点(用模拟 API保持向后兼容
_default_nodes = create_contact_nodes(_default_contact_api)
# 导出默认节点
parse_intent = _default_nodes["parse_intent"]
list_contacts = _default_nodes["list_contacts"]
add_contact = _default_nodes["add_contact"]
list_emails = _default_nodes["list_emails"]
generate_email_draft = _default_nodes["generate_email_draft"]
sniff_contacts = _default_nodes["sniff_contacts"]
format_result = _default_nodes["format_result"]
human_review = _default_nodes["human_review"]
send_email = _default_nodes["send_email"]
should_continue = _default_nodes["should_continue"]

View File

@@ -0,0 +1,104 @@
"""
通讯录子图状态定义
Contact Subgraph State Definition
"""
from enum import Enum, auto
from typing import Optional, Dict, List, Any
from dataclasses import dataclass, field
class ContactAction(Enum):
"""通讯录操作类型"""
NONE = auto()
CONTACT_LIST = auto() # 联系人列表
CONTACT_ADD = auto() # 添加联系人
CONTACT_UPDATE = auto() # 更新联系人
CONTACT_DELETE = auto() # 删除联系人
EMAIL_LIST = auto() # 邮件列表
EMAIL_READ = auto() # 读取邮件
EMAIL_SEND = auto() # 发送邮件
SNIFF_CONTACTS = auto() # 智能嗅探
@dataclass
class Contact:
"""联系人数据结构"""
id: Optional[str] = None
name: str = ""
phone: str = ""
email: str = ""
company: str = ""
position: str = ""
notes: str = ""
created_at: Optional[str] = None
updated_at: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class Email:
"""邮件数据结构"""
id: Optional[str] = None
subject: str = ""
sender: str = ""
recipients: List[str] = field(default_factory=list)
date: Optional[str] = None
body: str = ""
is_read: bool = False
mailbox: str = ""
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class ContactState:
"""通讯录子图状态"""
# ========== 输入 ==========
user_query: str = "" # 用户查询
user_id: str = "" # 用户ID
# 操作控制
action: ContactAction = ContactAction.NONE
action_params: Dict[str, Any] = field(default_factory=dict)
# ========== 执行过程 ==========
# 当前阶段
current_phase: str = "init" # init, processing, reviewing, done
# 联系人相关
contacts: List[Contact] = field(default_factory=list)
current_contact: Optional[Contact] = None
# 邮件相关
emails: List[Email] = field(default_factory=list)
current_email: Optional[Email] = None
# 邮件草稿(用于审核)
draft_subject: str = ""
draft_recipient: str = ""
draft_body: str = ""
# ========== 人工审核相关 ==========
pending_review: bool = False
review_type: str = "" # email_send, contact_delete
review_prompt: str = ""
review_approved: Optional[bool] = None
review_comment: str = ""
review_modified_content: str = ""
# ========== 智能嗅探 ==========
sniff_result: Optional[Dict[str, Any]] = None
sniffed_contacts: List[Contact] = field(default_factory=list)
sniff_confirmation_pending: bool = False
# ========== 结果 ==========
success: bool = False
error_message: str = ""
final_result: str = ""
result_data: Dict[str, Any] = field(default_factory=dict)
# ========== 元数据 ==========
start_time: Optional[str] = None
end_time: Optional[str] = None
duration: float = 0.0
debug_info: Dict[str, Any] = field(default_factory=dict)

View File

@@ -0,0 +1,428 @@
# 智能词典子图 (Dictionary Subgraph)
该子图负责处理翻译、查词、生词本管理等功能,基于 LangGraph 状态机编排多阶段学习流程支持联想记忆法、艾宾浩斯遗忘曲线复习、Anki 导出等核心能力。子图设计遵循"高效学习、科学复习、持久记忆"原则。
> **使用公共工具**意图理解、格式化输出、检查点持久化、条件路由、LLM 调用、数据库工具、状态基类
---
## 🎯 核心架构
### 技术栈
| 层级 | 组件 | 说明 |
|:-----|:-----|:-----|
| **编排框架** | LangGraph StateGraph | 状态机驱动的子图工作流编排 |
| **LLM 服务** | 智谱 AI / DeepSeek API | 翻译、释义生成、联想记忆、专业名词提炼(使用公共 LLM 工具) |
| **翻译服务** | DeepL / 有道 API | 高质量机器翻译 |
| **关系存储** | PostgreSQL | 生词本、复习记录持久化(使用公共数据库工具) |
| **导出工具** | csv / Anki APKG | 生词本导出格式 |
### 子图分层架构
```
┌─────────────────────────────────────────────────────────────────┐
│ 主图 (Main Graph) │
└──────────────────────────────┬──────────────────────────────────┘
│ 状态映射 / 结果聚合
┌─────────────────────────────────────────────────────────────────┐
│ 智能词典子图接口层 │
│ - 状态转换:主状态 ↔ 子图状态(使用公共状态基类) │
│ - 错误传播与优雅降级 │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 工作流编排层 │
│ - 节点调度与条件路由(使用公共路由工具) │
│ - 复习计划计算 │
│ - 状态持久化与检查点(使用公共检查点工具) │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 节点层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ │
│ │意图理解 │ │翻译节点 │ │查词节点 │ │每日一词 │ │专业提炼│ │
│ │(公共工具)│ │ │ │ │ │ │ │ │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ └────────┘ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │生词管理 │ │复习计划 │ │联想记忆 │ │格式输出 │ │
│ │ │ │ │ │ │ │(公共) │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
└──────────────────────────────┬──────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ 工具层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │翻译API │ │词典API │ │数据库工具│ │艾宾浩斯 │ │
│ │ │ │ │ │(公共) │ │ │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
└─────────────────────────────────────────────────────────────────┘
```
### 数据流总览
智能词典子图根据学习意图分支执行,支持查词、翻译、复习等多种学习模式。
```
用户请求
┌─────────────┐
│ 意图理解 │ ← 使用公共意图理解工具
└──────┬──────┘
├──────────┬──────────┬──────────┬──────────┬──────────┐
▼ ▼ ▼ ▼ ▼ ▼
翻译 查词 每日一词 专业提炼 生词管理 复习计划
│ │ │ │ │ │
▼ ▼ ▼ ▼ ▼ ▼
翻译API 词典API 每日推荐 术语提取 增删改查 复习计算
│ │ │ │ │ │
│ ▼ │ │ │ │
│ 联想记忆 │ │ │ │
│ │ │ │ │ │
└──────────┴──────────┴──────────┴──────────┴──────────┘
格式输出 ← 使用公共格式化工具
├──────────┐
│ │
▼ ▼
保存生词 Anki导出
│ │
└──────────┘
返回主图
```
---
## 📂 模块与文件结构
```
app/dictionary/
├── __init__.py
├── graph.py # 子图构建入口,定义状态图与路由
├── state.py # 子图状态定义(继承公共状态基类)
├── nodes/ # 节点实现
│ ├── __init__.py
│ ├── translate.py # 翻译节点
│ ├── lookup.py # 查词节点
│ ├── daily_word.py # 每日一词节点
│ ├── extract_terms.py # 专业名词提炼节点
│ ├── vocab.py # 生词本管理节点
│ ├── review.py # 复习计划节点
│ ├── association.py # 联想记忆节点
│ └── export.py # Anki导出节点
├── tools/ # 子图特有工具集
│ ├── translate_api.py # 翻译API工具
│ ├── dictionary_api.py # 词典API工具
│ ├── ebinghaus.py # 艾宾浩斯遗忘曲线工具
│ └── anki.py # Anki导出工具
└── persistence/ # (使用公共检查点工具,无需单独实现)
```
> **注意**:以下模块使用公共工具,无需单独实现:
> - 意图理解节点 → 使用 `agent_subgraphs.common.intent`
> - 格式输出节点 → 使用 `agent_subgraphs.common.format`
> - 检查点持久化 → 使用 `agent_subgraphs.common.checkpoint`
> - 条件路由 → 使用 `agent_subgraphs.common.routing`
> - LLM 调用 → 使用 `agent_subgraphs.common.llm`
> - 数据库操作 → 使用 `agent_subgraphs.common.db`
---
## 🎯 演进路线与核心机制
### Level 1基础翻译与查词
**核心机制**:调用翻译/词典 API展示基础释义。
- 支持中↔英、中↔日等多语言互译。
- 提供单词词性、释义、例句。
- 基础生词本功能(添加、查询)。
**适用场景**:快速翻译、单词查询。
**实现指引**:意图理解节点识别翻译/查词意图,路由到对应节点。
### Level 2联想记忆法
**核心机制**:利用 LLM 生成联想记忆法,帮助记忆单词。
- 词根词缀分析。
- 词源故事/文化背景。
- 趣味联想(谐音、画面感)。
- 同根词/近义词/反义词扩展。
**适用场景**:深度单词学习、高效记忆。
**实现指引**:查词后自动生成联想记忆内容,可选是否保存到生词本。
### Level 3艾宾浩斯遗忘曲线复习
**核心机制**:基于遗忘曲线科学安排复习时间。
- 首次学习后按 1天、2天、4天、7天、15天、30天 间隔复习。
- 根据用户记忆反馈动态调整复习间隔。
- 每日复习提醒。
**适用场景**:长期词汇积累、抗遗忘学习。
**实现指引**:复习计划节点计算下次复习时间,存入数据库。
### Level 4专业名词提炼与管理
**核心机制**:从文本中自动提取专业名词,建立术语库。
- 支持从任意文本中提取专业术语。
- 自动生成术语释义。
- 按领域分类管理术语库。
**适用场景**:专业文档阅读、学术学习。
**实现指引**:使用 LLM 进行 NER 和术语识别。
### Level 5智能词汇教练
**核心机制**:个性化学习路径、多模态记忆、学习进度追踪。
- 根据用户水平推荐学习内容。
- 图片、音频等多模态记忆辅助。
- 学习统计与进度可视化。
- 自适应难度调整。
**适用场景**:系统化语言学习、个性化辅导。
---
## 🔧 核心组件详解
### 1. 意图理解节点
**职责**:接收用户请求,区分翻译、查词、每日一词、专业提炼、生词管理、复习等意图。
**输入**:用户自然语言请求。
**输出**
- `intent_type`意图类别translate / lookup / daily / extract / vocab / review / export
- `target_word`:目标单词/文本。
- `source_lang`:源语言。
- `target_lang`:目标语言。
**实现要点**
- 使用 LLM 分类意图,输出结构化 JSON。
- 关键词匹配兜底(如"翻译"、"查一下")。
### 2. 翻译节点
**职责**:调用翻译 API返回高质量翻译结果。
**输入**:待翻译文本、源语言、目标语言。
**输出**
- `translation`:翻译结果。
- `alternative_translations`:备选翻译(如有)。
**实现要点**
- 优先使用 DeepL降级到有道或 LLM 翻译。
- 支持长文本分段翻译。
### 3. 查词节点
**职责**:查询单词详细释义、词性、例句等信息。
**输入**:目标单词。
**输出**
- `word_info`:单词信息(词性、释义、音标)。
- `examples`:例句列表。
**实现要点**
- 调用词典 API缺失时使用 LLM 生成。
- 支持英英、英汉双解。
### 4. 每日一词节点
**职责**:根据用户水平和历史,推荐今日学习单词。
**输入**:用户学习偏好。
**输出**
- `daily_word`:今日推荐单词。
- `word_detail`:单词详情。
- `learning_tip`:学习建议。
**实现要点**
- 结合用户历史生词和复习进度推荐。
- 难度适中递进。
### 5. 专业名词提炼节点
**职责**:从文本中提取专业名词,生成释义。
**输入**:待分析文本、领域(可选)。
**输出**
- `extracted_terms`:提取的专业名词列表。
- `term_definitions`:名词释义。
**实现要点**
- 使用 LLM 进行术语识别和定义生成。
- 支持按领域过滤。
### 6. 生词本管理节点
**职责**:生词本的增删改查操作。
**输入**:操作类型、生词数据。
**输出**
- `operation_result`:操作结果。
- `vocab_list`:更新后的生词列表。
**实现要点**
- 支持批量添加。
- 按标签/难度/复习时间筛选。
### 7. 复习计划节点
**职责**:基于艾宾浩斯遗忘曲线计算复习计划。
**输入**:生词 ID、上次复习时间、记忆强度。
**输出**
- `next_review`:下次复习时间。
- `review_schedule`:完整复习计划。
**实现要点**
- 使用标准艾宾浩斯间隔1天、2天、4天、7天、15天、30天
- 根据用户反馈动态调整。
### 8. 联想记忆节点
**职责**:为单词生成联想记忆法,帮助记忆。
**输入**:目标单词。
**输出**
- `root_analysis`:词根词缀分析。
- `etymology`:词源故事。
- `association`:趣味联想。
- `word_family`:同根词/近义词/反义词。
**实现要点**
- 使用 LLM 生成富有创意的记忆法。
- 支持用户自定义联想。
### 9. Anki 导出节点
**职责**:导出生词本为 Anki 可导入格式。
**输入**:导出范围(全部/按标签/按时间)。
**输出**
- `export_file`:导出文件路径。
- `export_count`:导出单词数量。
**实现要点**
- 支持 CSV 和 APKG 两种格式。
- 包含联想记忆内容。
---
## 🔀 条件路由详解
### 入口路由:意图分支
- **位置**:意图理解节点之后。
- **条件**
- `intent_type == "translate"` → 翻译节点。
- `intent_type == "lookup"` → 查词节点。
- `intent_type == "daily"` → 每日一词节点。
- `intent_type == "extract"` → 专业提炼节点。
- `intent_type == "vocab"` → 生词管理节点。
- `intent_type == "review"` → 复习计划节点。
- `intent_type == "export"` → Anki导出节点。
### 查词后续路由
- **位置**:查词节点之后。
- **条件**
- 用户询问"怎么记" → 联想记忆节点。
- 用户说"保存" → 生词管理节点(添加)。
- 无后续 → 格式输出。
---
## 📊 状态设计
### 状态结构概览
| 分组 | 字段 | 类型 | 说明 |
|:-----|:-----|:-----|:-----|
| **输入** | `user_input` | `str` | 用户原始请求 |
| **意图** | `intent_type` | `str` | 意图类别 |
| | `target_word` | `str` | 目标单词/文本 |
| | `source_lang` | `str` | 源语言 |
| | `target_lang` | `str` | 目标语言 |
| **翻译** | `translation` | `str` | 翻译结果 |
| | `alternative_translations` | `list[str]` | 备选翻译 |
| **查词** | `word_info` | `dict` | 单词信息 |
| | `examples` | `list[str]` | 例句 |
| **联想** | `root_analysis` | `str` | 词根分析 |
| | `etymology` | `str` | 词源 |
| | `association` | `str` | 联想记忆 |
| | `word_family` | `list[str]` | 词族 |
| **每日一词** | `daily_word` | `str` | 今日单词 |
| | `word_detail` | `dict` | 单词详情 |
| | `learning_tip` | `str` | 学习建议 |
| **专业提炼** | `extracted_terms` | `list[dict]` | 提取的术语 |
| | `term_definitions` | `dict` | 术语释义 |
| **生词本** | `vocab_list` | `list[dict]` | 生词列表 |
| | `operation_result` | `str` | 操作结果 |
| **复习** | `next_review` | `datetime` | 下次复习时间 |
| | `review_schedule` | `list[dict]` | 复习计划 |
| **导出** | `export_file` | `str` | 导出文件路径 |
| | `export_count` | `int` | 导出数量 |
| **控制流** | `current_phase` | `str` | 当前阶段 |
| | `next_node` | `str` | 下一节点 |
| **输出** | `final_result` | `str` | 最终结果 |
---
## 🔄 工作流程
### 查词+联想记忆流程
| 步骤 | 节点 | 说明 |
|:-----|:-----|:-----|
| 1 | 意图理解 | 识别查词意图 |
| 2 | 查词 | 查询单词释义 |
| 3 | 联想记忆 | 生成记忆法 |
| 4 | 询问保存 | 可选保存到生词本 |
| 5 | 格式输出 | 展示结果 |
### 复习流程
| 步骤 | 节点 | 说明 |
|:-----|:-----|:-----|
| 1 | 意图理解 | 识别复习意图 |
| 2 | 复习计划 | 获取今日需复习单词 |
| 3 | 复习交互 | 逐个复习,记录记忆强度 |
| 4 | 更新计划 | 计算下次复习时间 |
| 5 | 格式输出 | 展示复习结果 |
### Anki导出流程
| 步骤 | 节点 | 说明 |
|:-----|:-----|:-----|
| 1 | 意图理解 | 识别导出意图 |
| 2 | Anki导出 | 生成导出文件 |
| 3 | 格式输出 | 提供下载链接 |

View File

@@ -0,0 +1,50 @@
"""
词典子图 - 完善版
Dictionary Subgraph Module - Complete
"""
from .state import (
DictionaryState,
DictionaryAction,
WordEntry,
ExtractedTerm
)
from .graph import build_dictionary_subgraph
from .nodes import (
parse_intent,
query_word,
translate_text,
extract_terms,
get_daily_word,
lookup_word_book,
add_to_word_book,
format_result,
should_continue
)
from .api_client import dictionary_api, DictionaryAPIClient
__all__ = [
# State
"DictionaryState",
"DictionaryAction",
"WordEntry",
"ExtractedTerm",
# Graph
"build_dictionary_subgraph",
# Nodes
"parse_intent",
"query_word",
"translate_text",
"extract_terms",
"get_daily_word",
"lookup_word_book",
"add_to_word_book",
"format_result",
"should_continue",
# API
"dictionary_api",
"DictionaryAPIClient"
]

View File

@@ -0,0 +1,192 @@
"""
词典API调用工具
Dictionary API Client
支持 async 和真实数据库缓存
"""
from typing import Dict, Any, Optional
from dataclasses import dataclass
@dataclass
class DictionaryAPIClient:
"""
词典API客户端 - 可扩展支持多种API和数据库缓存
"""
# 可以配置多个API
youdao_api_key: Optional[str] = None
youdao_api_secret: Optional[str] = None
# 数据库 Repository可选用于缓存单词查询
word_repository: Optional[Any] = None
def __post_init__(self):
"""初始化后,如果有 repository 则支持 async"""
pass
async def query_word_db(self, user_id: str, word: str) -> Optional[Dict[str, Any]]:
"""从数据库缓存查询单词"""
if not self.word_repository:
return None
try:
entity = await self.word_repository.search_by_word(user_id, word)
if entity:
return {
"phonetic": entity.phonetic,
"part_of_speech": entity.part_of_speech,
"definitions": [entity.definition] if entity.definition else [],
"examples": [entity.examples] if entity.examples else []
}
except Exception as e:
print(f"从数据库查询单词失败:{e}")
return None
async def cache_word_db(self, user_id: str, word: str, data: Dict[str, Any]):
"""把单词查询结果缓存到数据库"""
if not self.word_repository:
return
try:
from ...db.models import WordEntity
entity = WordEntity(
user_id=user_id,
word=word,
phonetic=data.get("phonetic", ""),
part_of_speech=data.get("part_of_speech", ""),
definition=data.get("definitions", [""])[0] if data.get("definitions") else "",
examples=data.get("examples", [""])[0] if data.get("examples") else ""
)
await self.word_repository.insert(entity)
except Exception as e:
print(f"缓存单词到数据库失败:{e}")
async def query_word_youdao(self, word: str) -> Optional[Dict[str, Any]]:
"""
调用有道词典API查询单词async 版本)
注意需要配置有道API密钥才能使用
文档https://ai.youdao.com/doc.s#guide
"""
if not self.youdao_api_key or not self.youdao_api_secret:
return None
try:
# TODO: 实现真实的有道API调用用 httpx 或 aiohttp
# 这里是示例结构
return None
except Exception as e:
print(f"有道API调用失败{e}")
return None
async def translate_baidu(self, text: str, from_lang: str = "auto", to_lang: str = "zh") -> Optional[Dict[str, Any]]:
"""
调用百度翻译APIasync 版本)
注意需要配置百度API密钥才能使用
文档https://fanyi-api.baidu.com/doc/21
"""
# TODO: 实现真实的百度翻译API调用用 httpx 或 aiohttp
return None
def query_word_mock(self, word: str) -> Dict[str, Any]:
"""
模拟词典API - 目前用于演示
"""
mock_db = {
"serendipity": {
"phonetic": "/ˌserənˈdipədē/",
"part_of_speech": "n.",
"definitions": ["意外发现珍奇事物的能力", "机缘凑巧"],
"examples": ["Finding that old photo was pure serendipity."]
},
"ephemeral": {
"phonetic": "ˈfem(ə)rəl/",
"part_of_speech": "adj.",
"definitions": ["短暂的,瞬息的"],
"examples": ["Fame in the digital age is often ephemeral."]
},
"ubiquitous": {
"phonetic": "/yo͞oˈbikwədəs/",
"part_of_speech": "adj.",
"definitions": ["无处不在的", "普遍存在的"],
"examples": ["Smartphones have become ubiquitous in modern life."]
},
"eloquent": {
"phonetic": "/ˈeləkwənt/",
"part_of_speech": "adj.",
"definitions": ["雄辩的,有说服力的"],
"examples": ["She gave an eloquent speech at the conference."]
},
"resilient": {
"phonetic": "/rəˈzilyənt/",
"part_of_speech": "adj.",
"definitions": ["有复原力的,能适应的"],
"examples": ["The community has proven to be resilient in the face of challenges."]
}
}
if word.lower() in mock_db:
return mock_db[word.lower()]
else:
return {
"phonetic": "",
"part_of_speech": "n.",
"definitions": [f"{word}的释义1", f"{word}的释义2"],
"examples": [f"This is an example sentence with '{word}'."]
}
def translate_mock(self, text: str, from_lang: str = "auto", to_lang: str = "zh") -> Dict[str, Any]:
"""
模拟翻译API - 目前用于演示
"""
translations = {
"你好": "Hello",
"hello": "你好",
"人工智能": "Artificial Intelligence",
"artificial intelligence": "人工智能",
"ai": "人工智能",
"大模型": "Large Language Model",
"自然语言处理": "Natural Language Processing"
}
return {
"translated_text": translations.get(text.lower(), f"【翻译结果】{text}"),
"confidence": 0.95
}
def extract_terms_mock(self, text: str) -> list:
"""
模拟术语提取API
"""
return [
{"term": "AI", "type": "技术术语", "definition": "人工智能", "confidence": 0.95},
{"term": "LLM", "type": "技术术语", "definition": "大语言模型", "confidence": 0.92},
{"term": "NLP", "type": "技术术语", "definition": "自然语言处理", "confidence": 0.88}
]
# ========== 统一入口(优先查缓存) ==========
async def query_word(self, user_id: str = "default", word: str = "", use_cache: bool = True) -> Dict[str, Any]:
"""
查询单词(统一入口,优先查数据库缓存)
"""
# 1. 先查数据库缓存
if use_cache:
cached = await self.query_word_db(user_id, word)
if cached:
return cached
# 2. 查第三方 API暂未实现
api_result = await self.query_word_youdao(word)
if api_result:
if use_cache:
await self.cache_word_db(user_id, word, api_result)
return api_result
# 3. 用模拟数据(兜底)
mock_result = self.query_word_mock(word)
if use_cache:
await self.cache_word_db(user_id, word, mock_result)
return mock_result
# 单例实例(模拟模式,保持向后兼容)
dictionary_api = DictionaryAPIClient()

View File

@@ -0,0 +1,71 @@
"""
词典子图构建器 - 完善版
Dictionary Subgraph Builder - Complete
"""
from app.main_graph.graph import StateGraph, START, END
from .state import DictionaryState
from .nodes import (
parse_intent,
query_word,
translate_text,
extract_terms,
get_daily_word,
lookup_word_book,
add_to_word_book,
format_result,
should_continue
)
def build_dictionary_subgraph() -> StateGraph:
"""
构建词典子图
Returns:
配置好的 StateGraph
"""
# 创建图
graph = StateGraph(DictionaryState)
# 添加节点
graph.add_node("parse_intent", parse_intent)
graph.add_node("query_word", query_word)
graph.add_node("translate_text", translate_text)
graph.add_node("extract_terms", extract_terms)
graph.add_node("get_daily_word", get_daily_word)
graph.add_node("lookup_word_book", lookup_word_book)
graph.add_node("add_to_word_book", add_to_word_book)
graph.add_node("format_result", format_result)
# 添加边
# 从START开始
graph.add_edge(START, "parse_intent")
# 从parse_intent根据条件路由
graph.add_conditional_edges(
"parse_intent",
should_continue,
{
"query_word": "query_word",
"translate_text": "translate_text",
"extract_terms": "extract_terms",
"get_daily_word": "get_daily_word",
"lookup_word_book": "lookup_word_book",
"add_to_word_book": "add_to_word_book",
}
)
# 从各个操作节点到format_result
graph.add_edge("query_word", "format_result")
graph.add_edge("translate_text", "format_result")
graph.add_edge("extract_terms", "format_result")
graph.add_edge("get_daily_word", "format_result")
graph.add_edge("lookup_word_book", "format_result")
graph.add_edge("add_to_word_book", "format_result")
# 最终到END
graph.add_edge("format_result", END)
return graph

View File

@@ -0,0 +1,266 @@
"""
词典子图节点 - 使用公共工具版本
Dictionary Subgraph Nodes - Using Common Tools
"""
from typing import Dict, Any, List
from datetime import datetime
import random
# 公共工具
from ..common import (
MarkdownFormatter
)
from .state import (
DictionaryState,
DictionaryAction,
WordEntry,
ExtractedTerm
)
from .api_client import dictionary_api
# ========== 模拟生词本存储(后续可替换为数据库) ==========
WORD_BOOK_DB: Dict[str, List[Dict]] = {} # user_id -> [word_entries]
def parse_intent(state: DictionaryState) -> DictionaryState:
"""
解析用户意图节点(使用规则匹配)
确定用户想做什么操作
"""
# 子图特定的意图解析
query_lower = state.user_query.lower()
if any(keyword in query_lower for keyword in ["翻译", "translate", "英语", "英文"]):
state.action = DictionaryAction.TRANSLATE
state.action_params = {"text": state.user_query}
text = state.user_query
for keyword in ["翻译", "translate", "英语", "英文"]:
text = text.replace(keyword, "")
state.source_text = text.strip()
elif any(keyword in query_lower for keyword in ["查询", "query", "单词", "word"]):
state.action = DictionaryAction.QUERY
state.action_params = {"word": state.user_query}
elif any(keyword in query_lower for keyword in ["每日", "daily", "一词"]):
state.action = DictionaryAction.DAILY_WORD
elif any(keyword in query_lower for keyword in ["提取", "extract", "术语", "term"]):
state.action = DictionaryAction.EXTRACT
state.action_params = {"text": state.user_query}
else:
state.action = DictionaryAction.QUERY
state.action_params = {"word": state.user_query}
return state
def query_word(state: DictionaryState) -> DictionaryState:
"""
查询单词节点
"""
state.current_phase = "executing"
# 提取要查询的词
word = state.action_params.get("word", state.user_query)
# 清理关键词
for keyword in ["查询", "query", "单词", "word", "翻译", "translate", "英语", "英文"]:
word = word.replace(keyword, "").strip()
if not word:
word = "hello"
# 使用 API 客户端
word_entry = dictionary_api.query_word(word)
state.word_entry = word_entry
return state
def translate_text(state: DictionaryState) -> DictionaryState:
"""
翻译文本节点
"""
state.current_phase = "executing"
text = state.source_text or state.user_query
if not text:
# 清理关键词
for keyword in ["翻译", "translate", "英语", "英文"]:
text = text.replace(keyword, "").strip()
if not text:
text = "你好,世界!"
# 使用 API 客户端
translated = dictionary_api.translate(text)
state.source_text = text
state.translated_text = translated
return state
def extract_terms(state: DictionaryState) -> DictionaryState:
"""
提取术语节点
"""
state.current_phase = "executing"
text = state.action_params.get("text", state.user_query)
for keyword in ["提取", "extract", "术语", "term"]:
text = text.replace(keyword, "").strip()
if not text:
text = "Python is a great programming language for machine learning and data analysis."
# 使用 API 客户端
terms = dictionary_api.extract_terms(text)
state.extracted_terms = terms
return state
def get_daily_word(state: DictionaryState) -> DictionaryState:
"""
获取每日一词节点
"""
state.current_phase = "executing"
# 使用 API 客户端
word_entry = dictionary_api.get_daily_word()
state.daily_word = word_entry
return state
def lookup_word_book(state: DictionaryState) -> DictionaryState:
"""
查生词本节点
"""
state.current_phase = "executing"
user_id = state.user_id or "default"
if user_id not in WORD_BOOK_DB:
WORD_BOOK_DB[user_id] = []
state.word_book = WORD_BOOK_DB[user_id]
return state
def add_to_word_book(state: DictionaryState) -> DictionaryState:
"""
添加到生词本节点
"""
state.current_phase = "executing"
user_id = state.user_id or "default"
if user_id not in WORD_BOOK_DB:
WORD_BOOK_DB[user_id] = []
if state.word_entry:
entry_dict = {
"word": state.word_entry.word,
"phonetic": state.word_entry.phonetic,
"definition": state.word_entry.definition,
"added_at": datetime.now().isoformat()
}
WORD_BOOK_DB[user_id].append(entry_dict)
return state
def format_result(state: DictionaryState) -> DictionaryState:
"""
格式化结果节点(使用公共工具)
生成友好的 Markdown 输出
"""
state.current_phase = "formatting"
md = MarkdownFormatter()
output_lines = []
# 标题
output_lines.append("┌───────────────────────────────────┐")
output_lines.append("│ 📚 词典助手 │")
output_lines.append("└───────────────────────────────────┘")
output_lines.append("")
if state.action == DictionaryAction.QUERY and state.word_entry:
we = state.word_entry
output_lines.append(md.heading(f"📖 {we.word}", 2))
if we.phonetic:
output_lines.append(f"> 🔊 {we.phonetic}")
output_lines.append("")
output_lines.append(md.heading("释义", 3))
output_lines.append(md.bullet_list(we.definition))
if we.example_sentence:
output_lines.append("")
output_lines.append(md.heading("例句", 3))
output_lines.append(f"> {we.example_sentence}")
elif state.action == DictionaryAction.TRANSLATE and state.translated_text:
output_lines.append(md.heading("🌐 翻译结果", 2))
output_lines.append("")
output_lines.append(md.heading("原文", 3))
output_lines.append(f"> {state.source_text}")
output_lines.append("")
output_lines.append(md.heading("译文", 3))
output_lines.append(f"> {state.translated_text}")
elif state.action == DictionaryAction.EXTRACT and state.extracted_terms:
output_lines.append(md.heading("📝 提取的术语", 2))
output_lines.append("")
terms_data = [
{"术语": t.term, "释义": t.definition, "分类": t.category}
for t in state.extracted_terms
]
output_lines.append(md.table(terms_data))
elif state.action == DictionaryAction.DAILY_WORD and state.daily_word:
dw = state.daily_word
output_lines.append(md.heading("🌟 每日一词", 2))
output_lines.append("")
output_lines.append(md.heading(f"{dw.word}", 3))
if dw.phonetic:
output_lines.append(f"> 🔊 {dw.phonetic}")
output_lines.append("")
output_lines.append(md.bullet_list(dw.definition))
else:
output_lines.append(md.heading("✨ 操作完成", 2))
output_lines.append("您的请求已处理。")
# 页脚提示
output_lines.append("")
output_lines.append("---")
output_lines.append("💡 提示:您可以继续查询其他单词、翻译文本,或者提取术语!")
state.final_result = "\n".join(output_lines)
state.success = True
state.current_phase = "completed"
return state
def should_continue(state: DictionaryState) -> str:
"""
条件路由函数:根据 action 决定下一个节点
"""
action = state.action
if action == DictionaryAction.QUERY:
return "query_word"
elif action == DictionaryAction.TRANSLATE:
return "translate_text"
elif action == DictionaryAction.EXTRACT:
return "extract_terms"
elif action == DictionaryAction.DAILY_WORD:
return "get_daily_word"
elif action == DictionaryAction.LOOKUP_WORD_BOOK:
return "lookup_word_book"
elif action == DictionaryAction.ADD_TO_WORD_BOOK:
return "add_to_word_book"
else:
return "format_result"
return state

View File

@@ -0,0 +1,95 @@
"""
词典子图状态定义
Dictionary Subgraph State Definition
"""
from enum import Enum, auto
from typing import Optional, Dict, List, Any
from dataclasses import dataclass, field
class DictionaryAction(Enum):
"""词典操作类型"""
NONE = auto()
QUERY = auto() # 查询单词
TRANSLATE = auto() # 翻译文本
EXTRACT = auto() # 提取专业术语
DAILY_WORD = auto() # 每日一词
WORD_BOOK_LOOKUP = auto() # 生词本查询
WORD_BOOK_ADD = auto() # 添加到生词本
@dataclass
class WordEntry:
"""单词词条"""
word: str = ""
phonetic: str = "" # 音标
part_of_speech: str = "" # 词性
definitions: List[str] = field(default_factory=list) # 释义
examples: List[str] = field(default_factory=list) # 例句
synonyms: List[str] = field(default_factory=list) # 同义词
antonyms: List[str] = field(default_factory=list) # 反义词
source_language: str = "en" # 源语言
target_language: str = "zh" # 目标语言
in_word_book: bool = False # 是否在生词本
review_count: int = 0 # 复习次数
next_review_at: Optional[str] = None # 下次复习时间
created_at: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class ExtractedTerm:
"""提取的术语"""
term: str = ""
type: str = "" # 技术术语、医学术语等
definition: str = ""
context: str = ""
confidence: float = 0.0
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class DictionaryState:
"""词典子图状态"""
# ========== 输入 ==========
user_query: str = "" # 用户查询
user_id: str = "" # 用户ID
# 操作控制
action: DictionaryAction = DictionaryAction.NONE
action_params: Dict[str, Any] = field(default_factory=dict)
# 翻译专用
source_text: str = ""
source_language: str = "auto" # auto, en, zh, etc.
target_language: str = "zh" # 默认翻译成中文
# ========== 执行过程 ==========
current_phase: str = "init" # init, querying, extracting, done
# 查询结果
word_entry: Optional[WordEntry] = None
# 翻译结果
translated_text: str = ""
translation_confidence: float = 0.0
# 提取结果
extracted_terms: List[ExtractedTerm] = field(default_factory=list)
# 每日一词
daily_word: Optional[WordEntry] = None
daily_word_context: str = ""
# ========== 结果 ==========
success: bool = False
error_message: str = ""
final_result: str = ""
result_data: Dict[str, Any] = field(default_factory=dict)
# ========== 元数据 ==========
start_time: Optional[str] = None
end_time: Optional[str] = None
duration: float = 0.0
debug_info: Dict[str, Any] = field(default_factory=dict)

View File

@@ -0,0 +1,46 @@
"""
资讯子图 - 完善版
News Analysis Subgraph Module - Complete
"""
from .state import (
NewsAnalysisState,
NewsAction,
NewsItem,
NewsSource
)
from .graph import build_news_analysis_subgraph
from .nodes import (
parse_intent,
query_news,
analyze_url,
extract_keywords,
generate_report,
format_result,
should_continue
)
from .api_client import news_api, NewsAPIClient
__all__ = [
# State
"NewsAnalysisState",
"NewsAction",
"NewsItem",
"NewsSource",
# Graph
"build_news_analysis_subgraph",
# Nodes
"parse_intent",
"query_news",
"analyze_url",
"extract_keywords",
"generate_report",
"format_result",
"should_continue",
# API
"news_api",
"NewsAPIClient"
]

View File

@@ -0,0 +1,196 @@
"""
资讯子图API调用工具
News Analysis API Client
支持 async 和真实数据库缓存
"""
from typing import Dict, Any, Optional, List
import random
from datetime import datetime
from dataclasses import dataclass
@dataclass
class NewsAPIClient:
"""
资讯API客户端 - 可扩展支持多种API和数据库缓存
"""
# 可以配置多个API如 NewsAPI, 今日头条, 百度新闻等)
newsapi_key: Optional[str] = None
# 数据库 Repository可选用于缓存新闻
news_repository: Optional[Any] = None
async def query_news_db(self, user_id: str, keyword: str) -> Optional[List[Dict[str, Any]]]:
"""从数据库缓存查询新闻"""
if not self.news_repository:
return None
try:
entities = await self.news_repository.search_by_keywords(user_id, keyword)
if entities:
return [
{
"title": e.title,
"source": e.source,
"summary": e.content,
"keywords": e.keywords.split(",") if e.keywords else [],
"author": "",
"published_at": e.created_at
}
for e in entities
]
except Exception as e:
print(f"从数据库查询新闻失败:{e}")
return None
async def cache_news_db(self, user_id: str, news: Dict[str, Any]):
"""把新闻缓存到数据库"""
if not self.news_repository:
return
try:
from ...db.models import NewsEntity
entity = NewsEntity(
user_id=user_id,
title=news.get("title", ""),
content=news.get("summary", ""),
url=news.get("url", ""),
source=news.get("source", ""),
keywords=",".join(news.get("keywords", []))
)
await self.news_repository.insert(entity)
except Exception as e:
print(f"缓存新闻到数据库失败:{e}")
def query_news_mock(self, query: str) -> List[Dict[str, Any]]:
"""
模拟查询资讯 - 目前用于演示
"""
# 模拟资讯数据库
mock_news = [
{
"title": "OpenAI发布GPT-5智能再升级",
"source": "Tech News",
"summary": "最新消息OpenAI刚刚发布了GPT-5模型智能水平再次取得重大突破...",
"keywords": ["AI", "GPT-5", "OpenAI"],
"author": "AI Team",
"published_at": datetime.now().isoformat()
},
{
"title": "大模型在医疗领域的应用",
"source": "Health Tech",
"summary": "大模型AI技术正在医疗领域展现巨大潜力从辅助诊断到药物研发...",
"keywords": ["医疗", "大模型", "应用"],
"author": "Medical Team",
"published_at": datetime.now().isoformat()
},
{
"title": "2026年AI行业发展趋势报告",
"source": "Business Daily",
"summary": "最新行业报告显示AI行业将继续保持高速增长企业数字化转型加速...",
"keywords": ["趋势", "AI", "商业"],
"author": "Business Team",
"published_at": datetime.now().isoformat()
}
]
# 根据查询词简单过滤
results = []
query_lower = query.lower()
for news in mock_news:
if (query_lower in news["title"].lower() or
query_lower in news["summary"].lower() or
any(keyword.lower() in query_lower for keyword in news["keywords"])):
results.append(news)
# 如果没有匹配到,返回前两条
if not results:
results = mock_news[:2]
return results
def analyze_url_mock(self, url: str) -> Dict[str, Any]:
"""
模拟URL分析 - 目前用于演示
"""
return {
"title": f"分析结果:{url}",
"source": "URL Analyzer",
"summary": "已完成对该URL的内容分析包含文章摘要和情感倾向判断...",
"keywords": ["News", "Analysis", url.split("/")[-1] if url else "unknown"]
}
def extract_keywords_mock(self, text: str) -> List[str]:
"""
模拟关键词提取 - 目前用于演示
"""
# 简单的关键词提取模拟
common_keywords = ["AI", "大模型", "应用场景", "行业趋势", "创新", "技术"]
result = []
for keyword in common_keywords:
if keyword.lower() in text.lower():
result.append(keyword)
# 如果没找到,返回默认关键词
if not result:
result = ["AI", "大模型", "应用场景", "行业趋势"]
return result
def generate_report_mock(self, query: str) -> str:
"""
模拟报告生成 - 目前用于演示
"""
report = f"""═══════════════════════════════════════════
📊 资讯分析报告
═══════════════════════════════════════════
主题:{query}
📋 摘要:
这是一份关于 {query} 的资讯分析综合报告,包含最新行业动态和趋势分析。
🔍 主要发现:
1. AI技术持续快速发展
2. 大模型应用场景不断拓展
3. 行业数字化转型加速
🏷️ 关键词:
- AI
- 大模型
- 数字化转型
- 创新
═══════════════════════════════════════════
💡 建议:继续关注行业动态,把握发展机遇!
"""
return report
# ========== 统一入口(优先查缓存) ==========
async def query_news(self, user_id: str = "default", query: str = "", use_cache: bool = True) -> List[Dict[str, Any]]:
"""查询新闻(统一入口,优先查数据库缓存)"""
# 1. 先查数据库缓存
if use_cache:
cached = await self.query_news_db(user_id, query)
if cached:
return cached
# 2. 查第三方 API暂未实现
# api_result = await self.query_news_api(query)
# if api_result:
# for news in api_result:
# await self.cache_news_db(user_id, news)
# return api_result
# 3. 用模拟数据(兜底)
mock_result = self.query_news_mock(query)
if use_cache:
for news in mock_result:
await self.cache_news_db(user_id, news)
return mock_result
# 单例实例(模拟模式,保持向后兼容)
news_api = NewsAPIClient()

View File

@@ -0,0 +1,63 @@
"""
资讯子图构建器
News Analysis Subgraph Builder
"""
from app.main_graph.graph import StateGraph, START, END
from .state import NewsAnalysisState
from .nodes import (
parse_intent,
query_news,
analyze_url,
extract_keywords,
generate_report,
format_result,
should_continue
)
def build_news_analysis_subgraph() -> StateGraph:
"""
构建资讯子图
Returns:
配置好的 StateGraph
"""
# 创建图
graph = StateGraph(NewsAnalysisState)
# 添加节点
graph.add_node("parse_intent", parse_intent)
graph.add_node("query_news", query_news)
graph.add_node("analyze_url", analyze_url)
graph.add_node("extract_keywords", extract_keywords)
graph.add_node("generate_report", generate_report)
graph.add_node("format_result", format_result)
# 添加边
# 从START开始
graph.add_edge(START, "parse_intent")
# 从parse_intent根据条件路由
graph.add_conditional_edges(
"parse_intent",
should_continue,
{
"query_news": "query_news",
"analyze_url": "analyze_url",
"extract_keywords": "extract_keywords",
"generate_report": "generate_report",
}
)
# 从各个操作节点到format_result
graph.add_edge("query_news", "format_result")
graph.add_edge("analyze_url", "format_result")
graph.add_edge("extract_keywords", "format_result")
graph.add_edge("generate_report", "format_result")
# 最终到END
graph.add_edge("format_result", END)
return graph

View File

@@ -0,0 +1,185 @@
"""
资讯子图节点 - 使用公共工具版本
News Analysis Subgraph Nodes - Using Common Tools
"""
from typing import Dict, Any
from datetime import datetime
# 公共工具
from ..common import MarkdownFormatter
from .state import (
NewsAnalysisState,
NewsAction,
NewsItem,
NewsSource
)
from .api_client import news_api
def parse_intent(state: NewsAnalysisState) -> NewsAnalysisState:
"""
解析用户意图节点
确定用户想做什么操作
"""
query_lower = state.user_query.lower()
if any(keyword in query_lower for keyword in ["资讯", "新闻", "news", "report"]):
state.action = NewsAction.QUERY_NEWS
elif any(keyword in query_lower for keyword in ["分析", "analyze", "url", "链接"]):
state.action = NewsAction.ANALYZE_URL
elif any(keyword in query_lower for keyword in ["关键词", "keyword", "提取"]):
state.action = NewsAction.EXTRACT_KEYWORDS
elif any(keyword in query_lower for keyword in ["报告", "生成", "generate"]):
state.action = NewsAction.GENERATE_REPORT
else:
state.action = NewsAction.QUERY_NEWS
return state
def query_news(state: NewsAnalysisState) -> NewsAnalysisState:
"""
查询资讯节点
"""
state.current_phase = "executing"
# 使用 API 客户端
news_items = news_api.query_news(state.user_query)
state.news_items = news_items
return state
def analyze_url(state: NewsAnalysisState) -> NewsAnalysisState:
"""
分析 URL 节点
"""
state.current_phase = "executing"
# 从用户输入中提取 URL简单处理
query = state.user_query
url = query
for keyword in ["分析", "analyze", "url", "链接"]:
url = url.replace(keyword, "").strip()
if not url:
url = "https://example.com/news/article"
state.custom_urls = [url]
# 使用 API 客户端
analysis = news_api.analyze_url(url)
state.analysis = analysis
return state
def extract_keywords(state: NewsAnalysisState) -> NewsAnalysisState:
"""
提取关键词节点
"""
state.current_phase = "executing"
# 使用 API 客户端
keywords = news_api.extract_keywords(state.user_query)
state.extracted_keywords = keywords
return state
def generate_report(state: NewsAnalysisState) -> NewsAnalysisState:
"""
生成报告节点
"""
state.current_phase = "executing"
# 使用 API 客户端
report = news_api.generate_report(state.user_query)
state.report_content = report
return state
def format_result(state: NewsAnalysisState) -> NewsAnalysisState:
"""
格式化结果节点(使用公共工具)
"""
state.current_phase = "formatting"
md = MarkdownFormatter()
output_lines = []
output_lines.append("┌───────────────────────────────────┐")
output_lines.append("│ 📰 资讯助手 │")
output_lines.append("└───────────────────────────────────┘")
output_lines.append("")
if state.action == NewsAction.QUERY_NEWS and state.news_items:
output_lines.append(md.heading("📰 最新资讯", 2))
output_lines.append("")
for item in state.news_items:
output_lines.append(md.heading(item.title, 3))
output_lines.append(f"> 来源: {item.source.value}")
output_lines.append(f"> 时间: {item.published_at.strftime('%Y-%m-%d %H:%M')}")
if item.summary:
output_lines.append("")
output_lines.append(item.summary)
if item.url:
output_lines.append(f"🔗 链接: {md.link(item.title, item.url)}")
output_lines.append("")
elif state.action == NewsAction.EXTRACT_KEYWORDS and state.extracted_keywords:
output_lines.append(md.heading("🏷️ 提取的关键词", 2))
output_lines.append("")
keywords_data = [
{"关键词": k, "权重": f"{w:.2f}"}
for k, w in state.extracted_keywords.items()
]
output_lines.append(md.table(keywords_data))
elif state.action == NewsAction.GENERATE_REPORT and state.report_content:
output_lines.append(md.heading("📊 分析报告", 2))
output_lines.append("")
output_lines.append(state.report_content)
elif state.action == NewsAction.ANALYZE_URL and state.analysis:
output_lines.append(md.heading("🔍 URL 分析", 2))
output_lines.append("")
output_lines.append(f"> URL: {state.custom_urls[0]}")
output_lines.append("")
output_lines.append(state.analysis)
else:
output_lines.append(md.heading("✨ 操作完成", 2))
output_lines.append("您的请求已处理。")
# 页脚提示
output_lines.append("")
output_lines.append("---")
output_lines.append("💡 提示:您可以继续查询资讯、提取关键词或者生成报告!")
state.final_result = "\n".join(output_lines)
state.success = True
state.current_phase = "completed"
return state
def should_continue(state: NewsAnalysisState) -> str:
"""
条件路由函数:根据 action 决定下一个节点
"""
action = state.action
if action == NewsAction.QUERY:
return "query_news"
elif action == NewsAction.ANALYZE_URL:
return "analyze_url"
elif action == NewsAction.EXTRACT_KEYWORDS:
return "extract_keywords"
elif action == NewsAction.GENERATE_REPORT:
return "generate_report"
else:
return "format_result"

View File

@@ -0,0 +1,89 @@
"""
资讯子图状态定义
News Analysis Subgraph State Definition
"""
from enum import Enum, auto
from typing import Optional, Dict, List, Any
from dataclasses import dataclass, field
class NewsAction(Enum):
"""资讯操作类型"""
NONE = auto()
QUERY_NEWS = auto() # 查询资讯
ANALYZE_URL = auto() # 分析资讯
GENERATE_REPORT = auto() # 生成报告
FETCH_FROM_SOURCES = auto() # 从指定源获取
EXTRACT_KEYWORDS = auto() # 提取关键词
@dataclass
class NewsItem:
"""资讯条目"""
title: str = ""
url: str = ""
source: str = ""
content: str = ""
author: str = ""
published_at: Optional[str] = None
summary: str = ""
keywords: List[str] = field(default_factory=list)
sentiment: float = 0.0 # 情感分析得分
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class NewsSource:
"""资讯源"""
name: str = ""
url: str = ""
type: str = "" # rss, website, api
enabled: bool = True
last_fetched_at: Optional[str] = None
metadata: Dict[str, Any] = field(default_factory=dict)
@dataclass
class NewsAnalysisState:
"""资讯子图状态"""
# ========== 输入 ==========
user_query: str = "" # 用户查询
user_id: str = "" # 用户ID
# 操作控制
action: NewsAction = NewsAction.NONE
action_params: Dict[str, Any] = field(default_factory=dict)
# 源配置
use_follow_list: bool = False
custom_urls: List[str] = field(default_factory=list)
# ========== 执行过程 ==========
current_phase: str = "init" # init, fetching, analyzing, done
current_source_index: int = 0
primary_fetched: bool = False
# 源列表
sources: List[NewsSource] = field(default_factory=list)
# 资讯条目
news_items: List[NewsItem] = field(default_factory=list)
# 关键词
extracted_keywords: List[str] = field(default_factory=list)
# 报告
report_content: str = ""
# ========== 结果 ==========
success: bool = False
error_message: str = ""
final_result: str = ""
result_data: Dict[str, Any] = field(default_factory=dict)
# ========== 元数据 ==========
start_time: Optional[str] = None
end_time: Optional[str] = None
duration: float = 0.0
debug_info: Dict[str, Any] = field(default_factory=dict)