流式响应架构:AI Agent Harness 实时交互设计
流式响应架构:AI Agent Harness 实时交互设计
引言:为什么流式响应是 AI Agent 交互的未来
在当今快速发展的 AI 生态系统中,实时交互已经成为用户体验的核心要素。当我们与像 ChatGPT 这样的大型语言模型交互时,我们已经习惯了那种"打字机效果"——回答逐字逐句流式显示,而不是等待整个响应完成后一次性展示。这种看似简单的交互方式背后,隐藏着一套复杂的系统架构设计,我们称之为流式响应架构。
特别是在构建 AI Agent(智能代理)系统时,流式响应不仅仅是一个用户体验的增强,更是系统可靠性、效率和可扩展性的关键。在这篇文章中,我们将深入探讨如何设计和实现一个高效的 AI Agent Harness 实时交互系统,从底层原理到实际代码,从架构设计到最佳实践。
1. 核心概念:理解流式响应架构
1.1 什么是流式响应?
流式响应(Streaming Response)是一种数据传输模式,在这种模式下,服务器不会等待整个响应准备好后再发送,而是随着数据的生成逐步发送给客户端。这种方式与传统的"请求-响应"模式形成鲜明对比。
让我们通过一个简单的比喻来理解:想象你在餐厅点餐。在传统模式下,你点完菜后需要等待所有菜品准备好,然后服务员一次性把所有菜端上来。而在流式响应模式下,厨师每做好一道菜,服务员就立即端上来,你可以边吃边等下一道菜。
1.2 流式响应的核心优势
为什么流式响应在 AI Agent 交互中如此重要?主要有以下几个核心优势:
-
感知性能提升:用户可以立即看到响应的开始,而不是等待整个响应生成完毕。这种"即时反馈"大大提升了用户体验。
-
资源效率:服务器不需要在内存中保存整个响应,特别是对于长文本生成,这可以显著降低内存占用。
-
交互性增强:流式响应允许在生成过程中进行交互,例如用户可以在中途停止生成,或者系统可以在生成过程中插入中间结果。
-
错误恢复:如果连接中断,已经接收的部分响应仍然有用,用户不需要重新开始整个交互。
1.3 AI Agent Harness 的概念
在深入探讨架构之前,我们需要明确什么是 AI Agent Harness。简单来说,AI Agent Harness 是一个框架或平台,它提供了一套工具和基础设施,用于构建、部署和管理 AI 代理系统。
它就像是一个"智能代理的操作系统",处理诸如:
- 与各种 AI 模型的通信
- 状态管理
- 工具调用
- 内存管理
- 流式响应处理
- 错误恢复
等等核心功能。
2. 问题背景:从同步到流式的演变
2.1 传统同步交互的局限性
在早期的 AI 应用中,大多数交互都是同步的。用户发送一个请求,然后等待服务器处理完毕并返回完整响应。这种方式在处理简单、简短的请求时工作得很好,但随着 AI 模型能力的增强和应用场景的复杂化,它的局限性也日益明显:
-
用户体验差:对于需要较长时间生成的响应,用户会面对一个空白的屏幕或加载动画,不知道系统是否还在工作。
-
超时问题:HTTP 请求有超时限制,对于需要长时间处理的请求,容易导致连接中断。
-
资源浪费:服务器需要在内存中缓存整个响应,对于长文本生成,这可能导致内存压力。
-
缺乏交互性:用户无法在生成过程中提供反馈或中断操作。
让我们来看一个简单的同步交互的代码示例:
# 同步交互示例
import requests
def sync_query(question):
response = requests.post(
"https://api.example.com/query",
json={"question": question}
)
return response.json()["answer"]
# 使用
answer = sync_query("什么是量子计算?")
print(answer) # 等待整个响应返回后才打印
2.2 流式交互的兴起
随着 GPT 等大型语言模型的普及,流式交互逐渐成为标准。OpenAI 等公司率先在其 API 中提供了流式响应选项,使得开发者可以构建更具交互性的应用。
流式交互不仅仅是一个技术改进,它实际上改变了用户与 AI 系统交互的范式。用户不再是"发送请求,等待响应"的被动角色,而是可以参与到一个更加动态的对话过程中。
3. 核心架构设计:构建流式 AI Agent Harness
3.1 整体架构概览
一个完整的流式 AI Agent Harness 系统通常包含以下几个核心组件:
让我们逐一解析这些组件的功能:
- 客户端应用:用户直接交互的界面,可以是 Web 应用、移动应用或命令行工具。
- API 网关:处理客户端连接,负责协议转换、认证授权、限流等横切关注点。
- Agent 编排器:系统的核心,负责协调各个组件,管理对话流程,处理工具调用等。
- 工具注册表:管理 Agent 可以使用的各种工具,如网络搜索、代码执行等。
- 状态存储:保存对话状态、Agent 状态等重要信息。
- 模型网关:抽象不同 AI 模型提供商的接口,提供统一的流式调用能力。
- 事件总线:实现组件间的解耦通信,支持事件驱动架构。
- 可观测性系统:监控系统状态,收集日志、指标和追踪信息。
- 异步处理器:处理不需要实时返回的后台任务。
3.2 流式传输协议选择
在实现流式响应时,我们有几种主要的协议选择:
- WebSocket:全双工通信协议,允许服务器和客户端随时发送数据。
- Server-Sent Events (SSE):单向通信协议,只允许服务器向客户端推送数据。
- HTTP 分块传输编码:利用 HTTP/1.1 的分块传输功能实现流式响应。
让我们比较一下这三种协议:
| 特性 | WebSocket | SSE | 分块传输编码 |
|---|---|---|---|
| 双向通信 | ✅ | ❌ | ❌ |
| 浏览器原生支持 | ✅ | ✅ | ✅ |
| 自动重连 | ❌ | ✅ | ❌ |
| 事件类型支持 | 需要自定义 | ✅ | 需要自定义 |
| 防火墙友好性 | 可能被阻止 | ✅ | ✅ |
| 实现复杂度 | 中等 | 简单 | 简单 |
对于 AI Agent 交互场景,我们通常推荐使用 WebSocket 或 SSE,具体选择取决于是否需要双向通信。如果只需要服务器向客户端推送流式响应,SSE 是更简单的选择;如果需要在流式传输过程中客户端也能向服务器发送消息(如中断生成、提供反馈等),则 WebSocket 更合适。
4. 数学模型:流式响应的效率分析
4.1 感知延迟模型
让我们从数学角度分析流式响应如何改善用户体验。首先定义几个关键概念:
- TtotalT_{total}Ttotal:生成完整响应所需的总时间
- TfirstT_{first}Tfirst:第一个 token 到达客户端的时间
- TiT_{i}Ti:第 iii 个 token 到达的时间
- NNN:响应中的总 token 数
在传统的非流式响应中,用户需要等待 TtotalT_{total}Ttotal 时间才能看到任何内容。而在流式响应中,用户在 TfirstT_{first}Tfirst 时间就能开始看到内容,之后每隔 ΔTi=Ti−Ti−1\Delta T_i = T_i - T_{i-1}ΔTi=Ti−Ti−1 时间收到新的 token。
我们可以定义感知等待时间 WWW 为用户等待内容的主观感受。心理学研究表明,感知等待时间不仅仅是实际时间的函数,还与信息呈现的方式有关:
Wnon−streaming=TtotalW_{non-streaming} = T_{total}Wnon−streaming=Ttotal
Wstreaming=Tfirst+α⋅∑i=2NΔTiW_{streaming} = T_{first} + \alpha \cdot \sum_{i=2}^{N} \Delta T_iWstreaming=Tfirst+α⋅i=2∑NΔTi
其中 α\alphaα 是一个小于 1 的系数,表示一旦用户开始看到内容,后续等待的感知权重会降低。根据实际研究,α\alphaα 通常在 0.3 到 0.5 之间。
4.2 资源效率模型
让我们再分析一下资源效率。在非流式系统中,服务器需要在内存中保存完整的响应,直到它被完全发送。我们可以将内存使用量建模为:
Mnon−streaming(t)={k⋅∑i=1n(t)si,t≤Ttotal0,t>TtotalM_{non-streaming}(t) = \begin{cases} k \cdot \sum_{i=1}^{n(t)} s_i, & t \leq T_{total} \\ 0, & t > T_{total} \end{cases}Mnon−streaming(t)={k⋅∑i=1n(t)si,0,t≤Ttotalt>Ttotal
其中 n(t)n(t)n(t) 是到时间 ttt 为止生成的 token 数量,sis_isi 是第 iii 个 token 的大小,kkk 是存储开销系数。
而在流式系统中,一旦 token 被发送,就可以从内存中释放:
Mstreaming(t)=k⋅∑i=nsent(t)+1n(t)siM_{streaming}(t) = k \cdot \sum_{i=n_{sent}(t)+1}^{n(t)} s_iMstreaming(t)=k⋅i=nsent(t)+1∑n(t)si
其中 nsent(t)n_{sent}(t)nsent(t) 是到时间 ttt 为止已经发送的 token 数量。
假设生成和发送速率基本匹配,那么流式系统的内存使用量将大致保持在一个恒定的小值,而不是随着响应大小线性增长。
5. 核心算法:流式响应的生成与处理
5.1 Token 流式生成算法
在最底层,流式响应的核心是逐个生成并发送 token。让我们看一个简化的 Token 流式生成算法:
from typing import Generator, Any
import time
def token_stream_generator(
prompt: str,
model: Any,
delay_per_token: float = 0.05
) -> Generator[str, None, None]:
"""
模拟 AI 模型的 token 流式生成器
Args:
prompt: 用户输入的提示词
model: AI 模型实例
delay_per_token: 每个 token 生成的模拟延迟
Yields:
生成的 token
"""
# 模拟模型处理时间(第一个 token 前的延迟)
time.sleep(0.5)
# 模拟生成的响应
response = f"这是对'{prompt}'的流式响应。" \
"它会逐字逐句地显示出来," \
"就像你正在和一个真人对话一样。"
# 逐个生成 token(这里简化为逐个字符)
for char in response:
time.sleep(delay_per_token)
yield char
这个简单的生成器模拟了流式响应的基本行为:首先有一个初始延迟(模拟模型处理时间),然后逐个生成并 yield token。
5.2 带状态管理的流式处理算法
在实际的 AI Agent 系统中,我们不仅需要生成流式响应,还需要管理对话状态、处理工具调用等。让我们看一个更复杂的算法:
from typing import Generator, Dict, Any, Optional
from dataclasses import dataclass
from enum import Enum
import uuid
import time
class EventType(Enum):
"""流式事件类型"""
TEXT = "text"
TOOL_CALL = "tool_call"
TOOL_RESULT = "tool_result"
ERROR = "error"
DONE = "done"
@dataclass
class StreamEvent:
"""流式事件"""
event_type: EventType
content: Any
event_id: str = None
timestamp: float = None
def __post_init__(self):
if self.event_id is None:
self.event_id = str(uuid.uuid4())
if self.timestamp is None:
self.timestamp = time.time()
class AgentState:
"""Agent 状态管理"""
def __init__(self):
self.conversation_history = []
self.current_tool_calls = []
self.memory = {}
def add_message(self, role: str, content: str):
self.conversation_history.append({
"role": role,
"content": content
})
class AgentHarness:
"""Agent Harness 核心类"""
def __init__(self, model, tools=None):
self.model = model
self.tools = tools or {}
self.state = AgentState()
def process(self, user_input: str) -> Generator[StreamEvent, None, None]:
"""
处理用户输入,生成流式响应
Args:
user_input: 用户输入
Yields:
StreamEvent: 流式事件
"""
# 添加用户消息到状态
self.state.add_message("user", user_input)
try:
# 生成响应
yield from self._generate_response(user_input)
# 检查是否需要调用工具
if self._should_call_tool():
yield from self._handle_tool_calls()
except Exception as e:
yield StreamEvent(
event_type=EventType.ERROR,
content={"message": str(e)}
)
# 发送完成事件
yield StreamEvent(
event_type=EventType.DONE,
content={}
)
def _generate_response(self, user_input: str) -> Generator[StreamEvent, None, None]:
"""生成文本响应"""
# 模拟第一个 token 延迟
time.sleep(0.5)
# 这里简化处理,实际应该调用模型
response_text = f"我理解你的问题:'{user_input}'。让我思考一下..."
for char in response_text:
time.sleep(0.03) # 模拟生成延迟
yield StreamEvent(
event_type=EventType.TEXT,
content={"text": char}
)
def _should_call_tool(self) -> bool:
"""判断是否需要调用工具"""
# 简化逻辑,实际应该基于模型输出判断
return False
def _handle_tool_calls(self) -> Generator[StreamEvent, None, None]:
"""处理工具调用"""
# 工具调用逻辑,这里暂不实现
pass
这个算法展示了一个更完整的 Agent Harness 流式处理流程,包括状态管理、事件生成和错误处理。
6. 项目实战:构建一个简单的流式 AI Agent
6.1 项目介绍
让我们通过一个实际项目来应用前面讨论的概念。我们将构建一个名为 StreamAgent 的简单但功能完整的流式 AI Agent 系统,它具有以下特性:
- 基于 WebSocket 的双向流式通信
- 对话状态管理
- 模拟的工具调用能力
- 前端展示界面
6.2 技术栈选择
- 后端:Python + FastAPI (WebSocket 支持)
- 前端:HTML + JavaScript + CSS
- 状态管理:内存存储(简化版)
- AI 模型:模拟实现(方便演示)
6.3 环境安装
首先,让我们设置项目环境:
# 创建项目目录
mkdir stream-agent
cd stream-agent
# 创建虚拟环境
python -m venv venv
source venv/bin/activate # Windows 上使用 venv\Scripts\activate
# 安装依赖
pip install fastapi uvicorn python-multipart
6.4 系统架构设计
6.5 后端实现
让我们创建后端代码 main.py:
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from fastapi.staticfiles import StaticFiles
from fastapi.responses import FileResponse
from typing import Dict, Any, Generator, Optional
from dataclasses import dataclass
from enum import Enum
import uuid
import time
import json
import asyncio
app = FastAPI()
# 挂载静态文件
app.mount("/static", StaticFiles(directory="static"), name="static")
# 事件类型定义
class EventType(Enum):
TEXT = "text"
TOOL_CALL = "tool_call"
TOOL_RESULT = "tool_result"
ERROR = "error"
DONE = "done"
THINKING = "thinking"
@dataclass
class StreamEvent:
event_type: EventType
content: Any
event_id: str = None
timestamp: float = None
def __post_init__(self):
if self.event_id is None:
self.event_id = str(uuid.uuid4())
if self.timestamp is None:
self.timestamp = time.time()
def to_dict(self) -> Dict:
return {
"event_id": self.event_id,
"event_type": self.event_type.value,
"content": self.content,
"timestamp": self.timestamp
}
# 模拟工具
class ToolRegistry:
def __init__(self):
self.tools = {}
def register_tool(self, name: str, func, description: str):
self.tools[name] = {
"func": func,
"description": description
}
def execute_tool(self, name: str, params: Dict) -> Any:
if name not in self.tools:
raise ValueError(f"工具 '{name}' 不存在")
return self.tools[name]["func"](**params)
def get_tool_descriptions(self) -> str:
descriptions = []
for name, tool in self.tools.items():
descriptions.append(f"- {name}: {tool['description']}")
return "\n".join(descriptions)
# 模拟工具实现
def search_web(query: str) -> str:
"""模拟网络搜索工具"""
time.sleep(1) # 模拟搜索延迟
return f"[搜索结果] 关于 '{query}' 的信息:这是一个模拟的搜索结果,包含了一些相关信息。"
def calculate(expression: str) -> str:
"""模拟计算器工具"""
time.sleep(0.5) # 模拟计算延迟
try:
result = eval(expression)
return f"[计算结果] {expression} = {result}"
except Exception as e:
return f"[计算错误] {str(e)}"
# 初始化工具注册表
tool_registry = ToolRegistry()
tool_registry.register_tool("search_web", search_web, "搜索网络获取最新信息")
tool_registry.register_tool("calculate", calculate, "执行数学计算")
# Agent 状态管理
class AgentState:
def __init__(self, session_id: str):
self.session_id = session_id
self.conversation_history = []
self.memory = {}
def add_message(self, role: str, content: str):
self.conversation_history.append({
"role": role,
"content": content
})
def get_conversation_context(self) -> str:
return "\n".join([
f"{msg['role']}: {msg['content']}"
for msg in self.conversation_history
])
# 会话管理
class SessionManager:
def __init__(self):
self.sessions: Dict[str, AgentState] = {}
def create_session(self) -> str:
session_id = str(uuid.uuid4())
self.sessions[session_id] = AgentState(session_id)
return session_id
def get_session(self, session_id: str) -> Optional[AgentState]:
return self.sessions.get(session_id)
def end_session(self, session_id: str):
if session_id in self.sessions:
del self.sessions[session_id]
session_manager = SessionManager()
# StreamAgent 核心类
class StreamAgent:
def __init__(self, state: AgentState):
self.state = state
async def process(self, user_input: str) -> Generator[StreamEvent, None, None]:
"""处理用户输入,生成流式响应"""
# 添加用户消息到状态
self.state.add_message("user", user_input)
try:
# 思考阶段
yield StreamEvent(
event_type=EventType.THINKING,
content={"text": "正在思考..."}
)
await asyncio.sleep(0.5)
# 判断是否需要调用工具
if "搜索" in user_input or "查找" in user_input or "查询" in user_input:
# 需要搜索
tool_name = "search_web"
tool_params = {"query": user_input}
yield StreamEvent(
event_type=EventType.TOOL_CALL,
content={
"tool_name": tool_name,
"tool_params": tool_params
}
)
# 执行工具
tool_result = tool_registry.execute_tool(tool_name, tool_params)
yield StreamEvent(
event_type=EventType.TOOL_RESULT,
content={
"tool_name": tool_name,
"result": tool_result
}
)
# 添加工具结果到上下文
self.state.add_message("system", f"工具调用结果:{tool_result}")
# 基于工具结果生成回复
response_text = f"根据搜索结果,我找到了关于 '{user_input}' 的信息。让我为你总结一下:{tool_result}"
elif "计算" in user_input or "等于" in user_input or "+" in user_input or "-" in user_input or "*" in user_input or "/" in user_input:
# 需要计算
# 简单提取表达式
expression = user_input
for word in ["计算", "等于", "请问", "帮我", "一下"]:
expression = expression.replace(word, "")
expression = expression.strip()
tool_name = "calculate"
tool_params = {"expression": expression}
yield StreamEvent(
event_type=EventType.TOOL_CALL,
content={
"tool_name": tool_name,
"tool_params": tool_params
}
)
# 执行工具
tool_result = tool_registry.execute_tool(tool_name, tool_params)
yield StreamEvent(
event_type=EventType.TOOL_RESULT,
content={
"tool_name": tool_name,
"result": tool_result
}
)
# 生成回复
response_text = f"好的,我来帮你计算。{tool_result}"
else:
# 普通对话
response_texts = [
f"你说的是:'{user_input}'。这是一个很有趣的话题!",
"让我想想如何更好地回答这个问题...",
"在我看来,这个问题可以从多个角度来分析。",
"首先,我们需要理解问题的核心是什么。",
"然后,我们可以尝试寻找可能的解决方案。",
"当然,这只是我的个人观点,你可能有不同的看法。",
"你觉得这个方向对吗?我们可以进一步讨论。"
]
# 随机选择一个回复模板(实际应该由 AI 模型生成)
import random
response_text = random.choice(response_texts)
# 流式生成文本回复
for char in response_text:
await asyncio.sleep(0.03) # 模拟生成延迟
yield StreamEvent(
event_type=EventType.TEXT,
content={"text": char}
)
# 添加助手回复到状态
self.state.add_message("assistant", response_text)
except Exception as e:
yield StreamEvent(
event_type=EventType.ERROR,
content={"message": str(e)}
)
# 发送完成事件
yield StreamEvent(
event_type=EventType.DONE,
content={}
)
# WebSocket 连接管理
class ConnectionManager:
def __init__(self):
self.active_connections: Dict[str, WebSocket] = {}
async def connect(self, session_id: str, websocket: WebSocket):
await websocket.accept()
self.active_connections[session_id] = websocket
def disconnect(self, session_id: str):
if session_id in self.active_connections:
del self.active_connections[session_id]
async def send_event(self, session_id: str, event: StreamEvent):
if session_id in self.active_connections:
websocket = self.active_connections[session_id]
await websocket.send_json(event.to_dict())
manager = ConnectionManager()
# 路由
@app.get("/")
async def get():
return FileResponse("static/index.html")
@app.post("/api/session")
async def create_session():
session_id = session_manager.create_session()
return {"session_id": session_id}
@app.websocket("/ws/{session_id}")
async def websocket_endpoint(websocket: WebSocket, session_id: str):
await manager.connect(session_id, websocket)
session = session_manager.get_session(session_id)
if not session:
await websocket.close(code=1008, reason="无效的会话 ID")
return
agent = StreamAgent(session)
try:
while True:
data = await websocket.receive_json()
user_input = data.get("message")
if user_input:
async for event in agent.process(user_input):
await manager.send_event(session_id, event)
except WebSocketDisconnect:
manager.disconnect(session_id)
session_manager.end_session(session_id)
except Exception as e:
print(f"WebSocket 错误: {e}")
manager.disconnect(session_id)
session_manager.end_session(session_id)
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
6.6 前端实现
现在让我们创建前端界面。首先创建一个 static 目录,然后在其中创建 index.html:
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>StreamAgent - 流式 AI 代理演示</title>
<style>
* {
margin: 0;
padding: 0;
box-sizing: border-box;
}
body {
font-family: 'Segoe UI', Tahoma, Geneva, Verdana, sans-serif;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
height: 100vh;
display: flex;
justify-content: center;
align-items: center;
}
.container {
width: 100%;
max-width: 800px;
height: 90vh;
background: white;
border-radius: 16px;
box-shadow: 0 20px 60px rgba(0, 0, 0, 0.3);
display: flex;
flex-direction: column;
overflow: hidden;
}
.header {
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
padding: 20px;
text-align: center;
}
.header h1 {
font-size: 1.5rem;
margin-bottom: 5px;
}
.header p {
font-size: 0.9rem;
opacity: 0.9;
}
.chat-container {
flex: 1;
overflow-y: auto;
padding: 20px;
display: flex;
flex-direction: column;
gap: 15px;
}
.message {
max-width: 80%;
padding: 12px 16px;
border-radius: 16px;
line-height: 1.5;
animation: fadeIn 0.3s ease-in;
}
@keyframes fadeIn {
from { opacity: 0; transform: translateY(10px); }
to { opacity: 1; transform: translateY(0); }
}
.user-message {
align-self: flex-end;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
border-bottom-right-radius: 4px;
}
.assistant-message {
align-self: flex-start;
background: #f1f1f1;
color: #333;
border-bottom-left-radius: 4px;
}
.event-message {
align-self: center;
background: transparent;
color: #888;
font-size: 0.85rem;
font-style: italic;
padding: 5px 10px;
}
.tool-call {
background: #fff3cd;
color: #856404;
padding: 10px;
border-radius: 8px;
margin: 5px 0;
font-size: 0.9rem;
border: 1px solid #ffeaa7;
}
.tool-result {
background: #d4edda;
color: #155724;
padding: 10px;
border-radius: 8px;
margin: 5px 0;
font-size: 0.9rem;
border: 1px solid #c3e6cb;
}
.thinking-indicator {
display: inline-flex;
align-items: center;
gap: 5px;
}
.thinking-dots {
display: flex;
gap: 3px;
}
.thinking-dot {
width: 6px;
height: 6px;
background: #667eea;
border-radius: 50%;
animation: thinking 1.4s infinite ease-in-out both;
}
.thinking-dot:nth-child(1) { animation-delay: -0.32s; }
.thinking-dot:nth-child(2) { animation-delay: -0.16s; }
@keyframes thinking {
0%, 80%, 100% { transform: scale(0.8); opacity: 0.5; }
40% { transform: scale(1); opacity: 1; }
}
.input-container {
padding: 20px;
border-top: 1px solid #eee;
display: flex;
gap: 10px;
}
.input-container input {
flex: 1;
padding: 12px 16px;
border: 2px solid #eee;
border-radius: 24px;
font-size: 1rem;
transition: border-color 0.3s;
}
.input-container input:focus {
outline: none;
border-color: #667eea;
}
.input-container button {
padding: 12px 24px;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
border: none;
border-radius: 24px;
font-size: 1rem;
cursor: pointer;
transition: transform 0.2s, box-shadow 0.2s;
}
.input-container button:hover {
transform: translateY(-2px);
box-shadow: 0 4px 12px rgba(102, 126, 234, 0.4);
}
.input-container button:disabled {
opacity: 0.6;
cursor: not-allowed;
transform: none;
box-shadow: none;
}
.status-bar {
padding: 10px 20px;
background: #f8f9fa;
border-top: 1px solid #eee;
display: flex;
justify-content: space-between;
align-items: center;
font-size: 0.85rem;
color: #666;
}
.status-indicator {
display: inline-flex;
align-items: center;
gap: 5px;
}
.status-dot {
width: 8px;
height: 8px;
border-radius: 50%;
}
.status-dot.connected { background: #28a745; }
.status-dot.disconnected { background: #dc3545; }
.status-dot.connecting { background: #ffc107; animation: pulse 1s infinite; }
@keyframes pulse {
0%, 100% { opacity: 1; }
50% { opacity: 0.5; }
}
.welcome-message {
text-align: center;
padding: 40px 20px;
color: #666;
}
.welcome-message h2 {
margin-bottom: 10px;
color: #333;
}
.examples {
display: flex;
flex-wrap: wrap;
gap: 10px;
justify-content: center;
margin-top: 20px;
}
.example-btn {
padding: 8px 16px;
background: #f1f1f1;
border: none;
border-radius: 20px;
cursor: pointer;
font-size: 0.9rem;
transition: background 0.2s;
}
.example-btn:hover {
background: #e1e1e1;
}
</style>
</head>
<body>
<div class="container">
<div class="header">
<h1>StreamAgent</h1>
<p>流式 AI 代理交互演示</p>
</div>
<div class="chat-container" id="chatContainer">
<div class="welcome-message" id="welcomeMessage">
<h2>👋 你好!</h2>
<p>我是 StreamAgent,一个支持流式响应的 AI 代理。</p>
<p>试试下面的例子,开始我们的对话吧!</p>
<div class="examples">
<button class="example-btn" onclick="sendExample('什么是人工智能?')">什么是人工智能?</button>
<button class="example-btn" onclick="sendExample('搜索最新的科技新闻')">搜索最新的科技新闻</button>
<button class="example-btn" onclick="sendExample('计算 256 * 1024')">计算 256 * 1024</button>
</div>
</div>
</div>
<div class="input-container">
<input type="text" id="messageInput" placeholder="输入你的消息..." disabled>
<button id="sendButton" disabled>发送</button>
</div>
<div class="status-bar">
<div class="status-indicator">
<span class="status-dot connecting" id="statusDot"></span>
<span id="statusText">连接中...</span>
</div>
<div id="sessionInfo"></div>
</div>
</div>
<script>
let ws = null;
let sessionId = null;
let currentAssistantMessage = null;
const chatContainer = document.getElementById('chatContainer');
const messageInput = document.getElementById('messageInput');
const sendButton = document.getElementById('sendButton');
const statusDot = document.getElementById('statusDot');
const statusText = document.getElementById('statusText');
const sessionInfo = document.getElementById('sessionInfo');
const welcomeMessage = document.getElementById('welcomeMessage');
// 初始化
async function init() {
try {
// 创建会话
const response = await fetch('/api/session', { method: 'POST' });
const data = await response.json();
sessionId = data.session_id;
sessionInfo.textContent = `会话 ID: ${sessionId.substring(0, 8)}...`;
// 连接 WebSocket
connectWebSocket();
} catch (error) {
console.error('初始化失败:', error);
updateStatus('error', '初始化失败');
}
}
function connectWebSocket() {
updateStatus('connecting', '连接中...');
const wsUrl = `ws://${window.location.host}/ws/${sessionId}`;
ws = new WebSocket(wsUrl);
ws.onopen = () => {
updateStatus('connected', '已连接');
messageInput.disabled = false;
sendButton.disabled = false;
messageInput.focus();
};
ws.onmessage = (event) => {
const data = JSON.parse(event.data);
handleEvent(data);
};
ws.onclose = () => {
updateStatus('disconnected', '已断开连接');
messageInput.disabled = true;
sendButton.disabled = true;
// 尝试重连
setTimeout(() => {
if (ws && ws.readyState === WebSocket.CLOSED) {
connectWebSocket();
}
}, 3000);
};
ws.onerror = (error) => {
console.error('WebSocket 错误:', error);
updateStatus('error', '连接错误');
};
}
function updateStatus(status, text) {
statusText.textContent = text;
statusDot.className = 'status-dot';
switch (status) {
case 'connected':
statusDot.classList.add('connected');
break;
case 'disconnected':
statusDot.classList.add('disconnected');
break;
case 'connecting':
case 'error':
statusDot.classList.add('connecting');
break;
}
}
function handleEvent(data) {
const { event_type, content } = data;
switch (event_type) {
case 'thinking':
addThinkingMessage();
break;
case 'text':
appendToAssistantMessage(content.text);
break;
case 'tool_call':
addToolCallMessage(content.tool_name, content.tool_params);
break;
case 'tool_result':
addToolResultMessage(content.tool_name, content.result);
break;
case 'error':
addErrorMessage(content.message);
break;
case 'done':
finalizeAssistantMessage();
break;
}
// 滚动到底部
chatContainer.scrollTop = chatContainer.scrollHeight;
}
function addUserMessage(text) {
// 隐藏欢迎消息
welcomeMessage.style.display = 'none';
const messageDiv = document.createElement('div');
messageDiv.className = 'message user-message';
messageDiv.textContent = text;
chatContainer.appendChild(messageDiv);
}
function addThinkingMessage() {
const messageDiv = document.createElement('div');
messageDiv.className = 'message event-message';
messageDiv.innerHTML = `
<div class="thinking-indicator">
<span>正在思考</span>
<div class="thinking-dots">
<div class="thinking-dot"></div>
<div class="thinking-dot"></div>
<div class="thinking-dot"></div>
</div>
</div>
`;
chatContainer.appendChild(messageDiv);
}
function createAssistantMessage() {
currentAssistantMessage = document.createElement('div');
currentAssistantMessage.className = 'message assistant-message';
chatContainer.appendChild(currentAssistantMessage);
}
function appendToAssistantMessage(text) {
if (!currentAssistantMessage) {
createAssistantMessage();
}
currentAssistantMessage.textContent += text;
}
function finalizeAssistantMessage() {
currentAssistantMessage = null;
}
function addToolCallMessage(toolName, toolParams) {
const messageDiv = document.createElement('div');
messageDiv.className = 'message assistant-message';
messageDiv.innerHTML = `
<div class="tool-call">
<strong>🔧 调用工具:</strong> ${toolName}<br>
<strong>参数:</strong> ${JSON.stringify(toolParams)}
</div>
`;
chatContainer.appendChild(messageDiv);
}
function addToolResultMessage(toolName, result) {
const messageDiv = document.createElement('div');
messageDiv.className = 'message assistant-message';
messageDiv.innerHTML = `
<div class="tool-result">
<strong>✅ 工具结果 (${toolName}):</strong><br>
${result}
</div>
`;
chatContainer.appendChild(messageDiv);
}
function addErrorMessage(message) {
const messageDiv = document.createElement('div');
messageDiv.className = 'message event-message';
messageDiv.textContent = `❌ 错误: ${message}`;
chatContainer.appendChild(messageDiv);
}
function sendMessage() {
const message = messageInput.value.trim();
if (!message || !ws || ws.readyState !== WebSocket.OPEN) {
return;
}
addUserMessage(message);
ws.send(JSON.stringify({ message }));
messageInput.value = '';
}
function sendExample(example) {
messageInput.value = example;
sendMessage();
}
// 事件监听
sendButton.addEventListener('click', sendMessage);
messageInput.addEventListener('keypress', (e) => {
if (e.key === 'Enter') {
sendMessage();
}
});
// 启动
init();
</script>
</body>
</html>
6.7 运行项目
现在我们可以运行项目了:
# 创建 static 目录
mkdir static
# 确保 index.html 在 static 目录中
# 运行后端
python main.py
更多推荐
所有评论(0)