跳到主要内容

第三练 BOT

题目

  1. 利用LangChain、LangGraph等框架构建Agent,并尝试与常见聊天平台对接,作为前端交互界面(本例为飞书)
  2. 尝试亲自实现一个tools,例如早报模块,每天定时收集新闻,并发送到群里

项目解析

第2小问内容相对独立,因此单独开了一个news仓库

chat-agent/
├── main.py # 项目入口,启动飞书回调服务
├── pyproject.toml # 项目配置 & 依赖管理
├── README.md # 项目文档
├── .env / .env.example # 环境变量配置
├── .github/
│ └── PULL_REQUEST_TEMPLATE.md
├── app/
│ ├── __init__.py
│ ├── server.py # Flask OAuth 服务
│ ├── utils.py # 工具函数
│ │
│ ├── agents/ # Agent 模块
│ │ ├── __init__.py
│ │ ├── main_agent.py # 主 Agent(整合所有 tools)
│ │ ├── monitor.py # 群聊监控(主动回复)
│ │ ├── rag_workflow.py # RAG 工作流(LangGraph)
│ │ └── vlm.py # 视觉语言模型(图像理解)
│ │
│ ├── bot/ # 飞书机器人模块
│ │ ├── __init__.py
│ │ ├── client.py # 飞书 API 客户端
│ │ ├── auth.py # 认证相关
│ │ │
│ │ ├── calendar/
│ │ │ └── agenda.py # 日程管理工具
│ │ │
│ │ ├── contact/
│ │ │ ├── __init__.py
│ │ │ └── user.py # 用户信息获取
│ │ │
│ │ └── messages/
│ │ ├── __init__.py
│ │ ├── callback.py # 消息回调入口
│ │ ├── send.py # 消息发送
│ │ ├── get_history.py # 获取历史记录
│ │ └── get_msg_resource.py # 获取媒体资源
│ │
│ ├── tools/ # 工具模块
│ │ ├── __init__.py
│ │ └── web_search.py # 联网搜索
│ │
│ └── worker/ # 后台任务工作流
│ └── __init__.py # 长时间任务处理

├── storage/ # 文档存储(RAG 向量库)
│ ├── pdf/ # PDF 文档
│ └── txt/ # TXT 文档

└── tests/
└── test-lark.py # 飞书 SDK 测试

项目的核心代码都在app/里,一共分为这几个部分:agent/, bot/, tools/

1. agent/

agent/包含最主要的智能体功能。基础对话、工具调用、视觉、RAG都在这里实现。

main_agent.py是最基础的Agent,之后会负责飞书bot的基础对话和tool call的功能。

llm = ChatOpenAI(
...
)

system_prompt = f"""
现在是{datetime.datetime.now().strftime('%Y年%m月%d日 %H:%M:%S')}
你是一位优秀的助理,在飞书上和同事们交流,帮助同事们解决问题
...
"""

agent = create_agent(
llm, tools=[web_search, run_rag_workflow, comprehend_image, create_calendar_event_tool, ...], system_prompt=system_prompt
)


def generate_response(messages) -> str:
result = agent.invoke({"messages": messages})
return result["messages"][-1].content

你可能会发现,目前来看我们还没有维护上下文的操作。但是对于飞书机器人的项目来说其实不需要,因为飞书会帮我们维护上下文,也就是获取会话历史记录。这一部分会在后面bot/部分涉及。

对于私聊会话来说,对话逻辑是相对比较好设计的。在群聊中,通常情况下会设计成@机器人才触发它的回复。但既然有LLM的加持,我们也可以让模型自主决策在什么时候“插话”,例如当前会话中的用户遇到了问题,或者有一些重要信息的时候,模型此时应该出面解答,或者调用工具作为辅助。这时就需要monitor.py模块来随时监看群聊信息,设置结构化输出,让它判断何时说话

class IsNeeded(BaseModel):
is_needed: bool
query: str = Field(..., description="如果需要回答,请简要说明需要回答的内容;如果不需要回答,返回空字符串")

llm = ChatOpenAI(
model="qwen-flash",
max_retries=2,
api_key=os.getenv("DASHSCOPE_API_KEY"),
base_url=os.getenv("DASHSCOPE_API_BASE"),
)

llm_with_structure = llm.with_structured_output(IsNeeded)

def monitor(messages: list[dict]) -> IsNeeded:
...

为了让模型具备视觉,这里还放了一个vlm.py。不过为了保证主模型的性能,这里的视觉模型实际上是以subagent的形态呈现的

vlm = ChatOpenAI(
...
)

def encode_image(image_path):
...

@tool('comprehend_image', description="使用VLM模型理解图像内容。输入图像id和关于图像的问题,返回模型的回答。")
def comprehend_image(image: str, query: str) -> str:
'''使用VLM模型理解图像内容。
Args:
image: 图像id
query: 关于图像的问题
Returns:
模型的回答
'''
...

rag_workflow.py是RAG部分。考虑到现在的LLM上下文已经达到了百万级,因此,当文件大小不算太大的时候,直接调用一个低参数量的LLM进行RAG,省事且准确。如果需要处理的文件量非常大,再切换成常规embedding的方式

# 初始化模型
llm = ChatOpenAI(
...
)

# 初始化 embeddings
embeddings = DashScopeEmbeddings(
...
)


def estimate_tokens(text: str) -> int:
"""估算文本的token数量。
"""
...


def load_document(file_path: str) -> List[Document]:
"""加载文档。
"""
...


def get_total_tokens(file_paths: List[str]) -> int:
"""计算所有文件的总token数量。
"""
...


def llm_retriever_node(state: "RAGState") -> "RAGState":
"""使用LLM直接处理小文件的节点(总token < 990000)。
"""
...


def embedding_retriever_node(state: "RAGState") -> "RAGState":
"""使用embedding进行语义搜索的节点(适用于超大文件,总token >= 990000)。
"""
...

def route_based_on_size(state: "RAGState") -> str:
"""根据所有文件总token数路由到不同的处理节点。
"""
...


def combine_answers(state: "RAGState") -> "RAGState":
"""合并多个检索结果。
"""
# 目前使用单节点处理,后续可以扩展为多节点并行处理
return state


# 定义状态类型
class RAGState(TypedDict):
"""RAG工作流状态。"""
query: str
file_paths: List[str]
answer: str
source: str


def create_rag_workflow() -> StateGraph:
"""创建RAG工作流图。
"""
...

@tool("run_rag_workflow", description="基于RAG的文档问答工具,支持多文件输入。传入用户的查询和文件名列表,返回答案。")
def run_rag_workflow(query: str, file_names: List[str]) -> Dict[str, Any]:
"""运行RAG工作流。

Args:
query: 用户查询
file_names: 文件名列表(不需要输入完整路径,内部会自动拼接 ./storage/ 前缀)

Returns:
包含answer和source的字典
"""
...

2.bot/

这里主要是与飞书服务端API对接的部分,负责消息的接收、解析、发送。

2.1 消息回调 (messages/callback.py)

飞书机器人通过 WebSocket 长连接 接收消息事件。

def start_message_callback():
"""启动 WebSocket 长连接,监听消息事件"""

def do_p2_im_message_receive_v1(event):
"""飞书消息接收回调主入口"""

def process_message_async(message_data):
"""提交到线程池异步处理"""

消息处理流程:

接收消息事件 → 事件去重(非必须) → 线程池异步处理 → 消息筛选

私聊:直接处理
群聊+@机器人:直接处理
群聊+未@:触发 monitor 监控判断

调用 Agent 生成回复 → 发送消息

注意,对于用时较长的任务(例如Agent生成回答),使用常规的函数会阻塞进程,导致飞书服务端以为你没有收到消息,反复给你发当前的最新消息。因此处理函数一定要设计成异步的

2.2 消息发送 (messages/send.py)

和用户交互最关键的步骤。飞书的消息类型非常丰富,需要分别处理(例如text和post消息)。由于大模型的输出通常是带有markdown语法的,建议让机器人发送post消息

def send_text(chat_id: str, text: str):
"""发送纯文本消息"""

def send_post(chat_id: str, markdown: str, title: str = "回复"):
"""发送富文本消息(Markdown)"""

2.3 消息历史 (messages/get_history.py)

飞书会帮我们维护上下文,通过 get_history() 获取历史消息:

def get_history(chat_id: str) -> list[dict]:
"""获取会话历史消息"""

def history2dict(messages) -> list[dict]:
"""转换为 LangChain 格式"""

支持多种消息类型:text、post、image、file、audio 等,自动下载图片/文件到本地。

2.4 认证模块 (auth.py)

OAuth 2.0 认证流程,这一部分很重要,飞书的日历日程管理、获取用户信息等功能需要获取user_access_token,否则没有权限。

def get_authorization_url() -> str:
"""生成授权 URL"""

def get_tenant_access_token() -> str:
"""获取租户 access_token(应用级别)"""

...

2.5 日程管理 (calendar/agenda.py)

日程管理的API可传输的参数量很大,建议只保留常用的即可,并且针对不同的参数做好数据结构的设计。

def create_calendar_event(config: CalendarEventConfig):
"""创建飞书日程事件"""

@tool
def create_calendar_event_tool():
"""封装为 LangChain Tool,供 Agent 调用"""

2.6 用户信息 (contact/user.py)

可以基于此查看到会话中用户的相关信息,例如姓名,这样Agent就能看到带着名字的上下文而不是open_id,更能分清对话角色了

def get_users(user_ids: list[str]) -> list[dict]:
"""批量获取用户信息"""

3.tools/

这里是 Agent 会用到的外部工具。

3.1 联网搜索 (tools/web_search.py)

这里的示例使用的是Tavily API。也可以更换成其他的网络搜索的API。但是要注意完善prompt engineering,要求Agent在消息的最后注明引用来源,或者在文中以markdown超链接的语法标注。

@tool("web_search")
def web_search(query: str, max_results: int = 5, ...):
"""使用 Tavily 进行网络搜索"""

关键技术点

  1. LangChain Agent: 整合多个工具,统一工具调用入口
  2. LangGraph RAG: 根据文档大小智能路由,小文档直接处理,大文档向量检索
  3. 群聊监控: 未@机器人时,通过 LLM 判断是否需要主动介入
  4. WebSocket 长连接: 实时接收飞书消息事件
  5. 线程池异步处理: 避免阻塞消息回调,影响飞书服务