引言:AI应用开发的新范式

在AI技术飞速发展的今天,如何快速构建一个功能完善、用户体验优秀的AI聊天应用?AG-UI协议应运而生!本文将带你从零开始,完整搭建一个基于AG-UI协议的智能聊天系统,包括本地AI模型部署、后端API开发和前端交互实现。

完整前端代码:https://blog.csdn.net/weixin_40582941/article/details/155032876?spm=1001.2014.3001.5501

第一章:什么是AG-UI?AI应用开发的"通用语言"

AG-UI协议的核心价值

AG-UI(Agent-User Interface)是一个专门为AI智能体交互设计的标准化通信协议。想象一下,如果每个AI应用都需要重新设计通信方式,那将是多么低效!AG-UI就像AI世界的"HTTP协议",为开发者提供了一套统一的交互标准。

为什么选择AG-UI?

  • 🎯 标准化事件系统:预定义完整的事件类型(开始、内容、完成、错误)

  • ⚡ 实时流式传输:支持逐词输出的流式响应,告别等待焦虑

  • 🔄 多会话管理:内置会话管理,支持多用户同时聊天

  • 🔌 前后端解耦:前端无需关心后端使用什么AI服务

  • 🚀 快速集成:几分钟内接入任何AI模型服务

AG-UI事件流详解

第二章:环境准备 - 搭建你的AI基础设施

2.1 部署本地Ollama服务

步骤1:安装Ollama

Windows系统:

bash

# 方法1:直接下载安装包
访问 https://ollama.ai/download 下载安装程序

# 方法2:使用winget(Windows 11推荐)
winget install Ollama.Ollama

macOS系统:

bash

# 方法1:直接下载安装包
访问 https://ollama.ai/download 下载DMG文件

# 方法2:使用Homebrew
brew install ollama

Linux系统:

bash

# Ubuntu/Debian
curl -fsSL https://ollama.ai/install.sh | sh

# CentOS/RHEL/Fedora
curl -fsSL https://ollama.ai/install.sh | sh
步骤2:启动Ollama服务

Windows:

bash

# 安装后Ollama会自动启动服务
# 手动启动方式:
ollama serve

macOS/Linux:

bash

# 启动服务
ollama serve

# 或者使用systemd(Linux)
sudo systemctl enable ollama
sudo systemctl start ollama
步骤3:下载AI模型

bash

# 下载轻量级模型(推荐初学者)
ollama pull gemma2:2b

# 或者下载其他流行模型
ollama pull llama2:7b
ollama pull mistral:7b
ollama pull qwen2:1.5b

# 查看已安装的模型
ollama list
步骤4:验证安装

bash

# 测试模型响应
ollama run gemma2:2b "Hello, how are you?"

# 检查服务状态
curl http://localhost:11434/api/tags

2.2 Python环境配置

bash

# 创建虚拟环境(推荐)
python -m venv agui_env
source agui_env/bin/activate  # Linux/macOS
# agui_env\Scripts\activate  # Windows

# 安装必要依赖
pip install fastapi uvicorn websockets requests

# 安装AG-UI协议库
pip install ag-ui-protocol

第三章:后端开发 - 构建AG-UI智能聊天API

3.1 完整后端代码实现

python

from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Request
from fastapi.responses import StreamingResponse
from fastapi.middleware.cors import CORSMiddleware
from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates
import asyncio
import uuid
from datetime import datetime
from typing import Dict, List, Optional
import json
import requests
import os

# AG-UI Protocol Components
from ag_ui.encoder import EventEncoder
from ag_ui.core import (
    RunStartedEvent,
    RunFinishedEvent, 
    TextMessageContentEvent,
    RunErrorEvent
)
from ag_ui.core import (
    EventType,
)

# 初始化FastAPI应用
app = FastAPI(
    title="AG-UI IM Message Demo",
    description="基于AG-UI协议的智能聊天应用后端",
    version="1.0.0"
)

# 配置CORS中间件,允许前端跨域访问
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],  # 生产环境应设置为具体的前端地址
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

# 初始化AG-UI组件
encoder = EventEncoder()

# 内存中存储聊天会话(生产环境建议使用Redis或数据库)
chat_sessions: Dict[str, List[Dict]] = {}

class ConnectionManager:
    """WebSocket连接管理器"""
    def __init__(self):
        self.active_connections: List[WebSocket] = []

    async def connect(self, websocket: WebSocket):
        await websocket.accept()
        self.active_connections.append(websocket)

    def disconnect(self, websocket: WebSocket):
        if websocket in self.active_connections:
            self.active_connections.remove(websocket)

    async def send_personal_message(self, message: str, websocket: WebSocket):
        await websocket.send_text(message)

manager = ConnectionManager()

def call_ollama_service(message: str) -> List[str]:
    """
    调用本地Ollama服务并返回响应tokens
    
    Args:
        message: 用户输入的消息
        
    Returns:
        List[str]: 分词后的响应内容
    """
    try:
        # Ollama API端点(默认本地11434端口)
        url = "http://localhost:11434/api/generate"
        
        # 准备请求数据
        payload = {
            "model": "gemma2:2b",  # 替换为你下载的模型名称
            "prompt": message,
            "stream": False,       # 设为True可启用原生流式传输
            "options": {
                "temperature": 0.7,   # 控制创造性(0-1)
                "top_p": 0.9,        # 核采样参数
                "top_k": 40,         # 顶级K采样
            }
        }
        
        # 调用Ollama API
        response = requests.post(url, json=payload, timeout=60)
        response.raise_for_status()
        
        # 解析响应
        result = response.json()
        full_response = result.get("response", "")
        
        # 将完整响应拆分为tokens用于流式模拟
        # 在实际项目中,你可以使用Ollama的原生流式API
        tokens = []
        words = full_response.split()
        for word in words:
            tokens.append(word + " ")
        
        return tokens if tokens else ["I received your message but got no response from Ollama."]
        
    except requests.exceptions.ConnectionError:
        return ["❌ 错误:无法连接到Ollama服务。请确保Ollama正在localhost:11434上运行"]
    except requests.exceptions.Timeout:
        return ["⏰ 错误:Ollama服务响应超时。请求时间过长"]
    except requests.exceptions.RequestException as e:
        return [f"🔧 错误调用Ollama服务:{str(e)}"]
    except Exception as e:
        return [f"⚠️ 意外错误:{str(e)}"]

# 备用AI响应函数(当Ollama不可用时使用)
def simulate_ai_thinking(message: str) -> List[str]:
    """模拟AI处理消息并返回响应tokens"""
    responses = {
        "hello": ["你好!", "我是你的AI助手。", "今天有什么可以帮你的吗?"],
        "天气": ["今天天气晴朗", "最高温度25℃。", "非常适合外出活动!"],
        "时间": [f"当前时间是 {datetime.now().strftime('%H:%M')}。", "祝你今天愉快!"]
    }
    
    msg_lower = message.lower()
    for key, response in responses.items():
        if key in msg_lower:
            return response
    
    return [f"我收到你的消息:'{message}'。", "这是一个基于AG-UI协议的模拟响应。"]

@app.websocket("/ws/chat/{session_id}")
async def websocket_chat(websocket: WebSocket, session_id: str):
    """
    WebSocket聊天端点 - 支持实时双向通信
    
    Args:
        websocket: WebSocket连接对象
        session_id: 会话ID
    """
    await manager.connect(websocket)
    
    # 初始化会话(如果不存在)
    if session_id not in chat_sessions:
        chat_sessions[session_id] = []
    
    try:
        while True:
            # 接收客户端消息
            data = await websocket.receive_text()
            
            try:
                # 解析传入消息
                message_data = json.loads(data)
                user_message = message_data.get("message", "")
                
                if not user_message:
                    # 发送错误事件
                    error_event = RunErrorEvent(
                        type=EventType.ERROR,
                        message="消息不能为空"
                    )
                    await websocket.send_text(encoder.encode(error_event))
                    continue
                
                # 存储用户消息到会话历史
                chat_sessions[session_id].append({
                    "role": "user",
                    "message": user_message,
                    "timestamp": datetime.now().isoformat()
                })
                
                # 🔄 AG-UI协议事件流开始
                
                # 1. 发送RUN_STARTED事件
                run_id = f"run_{uuid.uuid4().hex}"
                thread_id = session_id  # 使用session_id作为thread_id
                
                run_started = RunStartedEvent(
                    type=EventType.RUN_STARTED,
                    run_id=run_id,
                    thread_id=thread_id
                )
                await websocket.send_text(encoder.encode(run_started))
                
                # 2. 处理消息并流式传输响应
                try:
                    response_tokens = call_ollama_service(user_message)
                except:
                    # 如果Ollama服务不可用,使用模拟响应
                    response_tokens = simulate_ai_thinking(user_message)
                
                # 流式传输响应tokens(AG-UI TextMessageContentEvent)
                full_response = ""
                for token in response_tokens:
                    await asyncio.sleep(0.05)  # 模拟思考时间,创造更自然的体验
                    full_response += token
                    
                    text_event = TextMessageContentEvent(
                        type=EventType.TEXT_MESSAGE_CONTENT,
                        message_id=f"msg_{uuid.uuid4().hex}",
                        thread_id=thread_id,
                        run_id=run_id,
                        delta=token  # 增量内容
                    )
                    await websocket.send_text(encoder.encode(text_event))
                
                # 3. 发送RUN_FINISHED事件
                run_finished = RunFinishedEvent(
                    type=EventType.RUN_FINISHED,
                    run_id=run_id,
                    thread_id=thread_id
                )
                await websocket.send_text(encoder.encode(run_finished))
                
                # 存储AI响应到会话历史
                chat_sessions[session_id].append({
                    "role": "assistant",
                    "message": full_response,
                    "timestamp": datetime.now().isoformat()
                })
                
            except json.JSONDecodeError:
                # JSON解析错误处理
                error_event = RunErrorEvent(
                    type=EventType.ERROR,
                    thread_id=session_id,
                    error="无效的JSON格式"
                )
                await websocket.send_text(encoder.encode(error_event))
            except Exception as e:
                # 其他错误处理
                error_event = RunErrorEvent(
                    type=EventType.ERROR,
                    thread_id=session_id,
                    error=f"处理错误:{str(e)}"
                )
                await websocket.send_text(encoder.encode(error_event))
                
    except WebSocketDisconnect:
        # WebSocket连接断开处理
        manager.disconnect(websocket)
        print(f"WebSocket连接断开: {session_id}")

@app.post("/api/chat/{session_id}/stream")
async def http_stream_chat(session_id: str, request: Request):
    """
    HTTP流式聊天端点 - 支持Server-Sent Events (SSE)
    
    适用于不支持WebSocket的环境或需要HTTP-only的场景
    """
    try:
        data = await request.json()
        user_message = data.get("message", "")
        
        if not user_message:
            async def error_generator():
                error_event = RunErrorEvent(
                    type=EventType.ERROR,
                    thread_id=session_id,
                    error="消息不能为空"
                )
                yield encoder.encode(error_event)
            return StreamingResponse(error_generator(), media_type="text/plain")
        
        # 初始化会话
        if session_id not in chat_sessions:
            chat_sessions[session_id] = []
        
        # 存储用户消息
        chat_sessions[session_id].append({
            "role": "user",
            "message": user_message,
            "timestamp": datetime.now().isoformat()
        })
        
        async def event_generator():
            """AG-UI事件生成器"""
            try:
                # 1. 开始运行
                run_id = f"run_{uuid.uuid4().hex}"
                thread_id = session_id
                
                run_started = RunStartedEvent(
                    type=EventType.RUN_STARTED,
                    run_id=run_id,
                    thread_id=thread_id
                )
                yield encoder.encode(run_started)
                
                # 2. 处理和流式传输响应
                try:
                    response_tokens = call_ollama_service(user_message)
                except:
                    response_tokens = simulate_ai_thinking(user_message)
                
                for token in response_tokens:
                    await asyncio.sleep(0.05)
                    text_event = TextMessageContentEvent(
                        type=EventType.TEXT_MESSAGE_CONTENT,
                        message_id=f"msg_{uuid.uuid4().hex}",
                        delta=token
                    )
                    yield encoder.encode(text_event)
                
                # 3. 完成运行
                run_finished = RunFinishedEvent(
                    type=EventType.RUN_FINISHED,
                    run_id=run_id,
                    thread_id=thread_id
                )
                yield encoder.encode(run_finished)
                
                # 存储完整响应
                full_response = "".join(response_tokens)
                chat_sessions[session_id].append({
                    "role": "assistant",
                    "message": full_response,
                    "timestamp": datetime.now().isoformat()
                })
                
            except Exception as e:
                # 错误处理
                error_event = RunErrorEvent(
                    message=f"流式传输错误:{str(e)}"
                )
                yield encoder.encode(error_event)
        
        return StreamingResponse(
            event_generator(), 
            media_type="text/plain",
            headers={
                "Cache-Control": "no-cache",
                "Connection": "keep-alive",
                "Access-Control-Allow-Origin": "*",
            }
        )
        
    except Exception as e:
        # 返回HTTP错误响应
        from fastapi import HTTPException
        raise HTTPException(status_code=500, detail=f"服务器错误:{str(e)}")

# 会话管理API
@app.get("/api/sessions/{session_id}/messages")
async def get_chat_history(session_id: str):
    """获取指定会话的聊天历史"""
    if session_id not in chat_sessions:
        return {"session_id": session_id, "messages": []}
    
    return {
        "session_id": session_id,
        "messages": chat_sessions[session_id]
    }

@app.post("/api/sessions/{session_id}/clear")
async def clear_chat_history(session_id: str):
    """清空指定会话的聊天历史"""
    if session_id in chat_sessions:
        chat_sessions[session_id] = []
    return {"status": "cleared", "session_id": session_id}

@app.get("/api/sessions/new")
async def create_new_session():
    """创建新的聊天会话"""
    session_id = str(uuid.uuid4())
    chat_sessions[session_id] = []
    return {"session_id": session_id}

# 系统状态检查端点
@app.get("/")
async def root():
    return {
        "message": "AG-UI IM Message Demo API", 
        "status": "running",
        "timestamp": datetime.now().isoformat(),
        "active_sessions": len(chat_sessions)
    }

@app.get("/health")
async def health_check():
    """健康检查端点"""
    return {
        "status": "healthy", 
        "timestamp": datetime.now().isoformat(),
        "ollama_status": "unknown"  # 可以添加Ollama服务状态检查
    }

@app.get("/api/models")
async def get_available_models():
    """获取可用的Ollama模型列表"""
    try:
        response = requests.get("http://localhost:11434/api/tags", timeout=5)
        if response.status_code == 200:
            return response.json()
        else:
            return {"error": "无法获取模型列表"}
    except:
        return {"error": "Ollama服务未运行"}

# 启动应用
if __name__ == "__main__":
    import uvicorn
    
    print("🚀 启动AG-UI聊天服务器...")
    print("📍 API文档: http://localhost:8000/docs")
    print("🔗 健康检查: http://localhost:8000/health")
    print("🤖 模型列表: http://localhost:8000/api/models")
    
    uvicorn.run(
        app, 
        host="0.0.0.0",  # 允许外部访问
        port=8000, 
        log_level="info",
        reload=True  # 开发时启用热重载
    )

3.2 后端代码深度解析

AG-UI事件流的精妙设计

python

# AG-UI事件序列示例
events_sequence = [
    "RUN_STARTED",           # 开始处理
    "TEXT_MESSAGE_CONTENT",  # 流式内容(多次)
    "TEXT_MESSAGE_CONTENT",  # 流式内容(多次)
    "RUN_FINISHED",          # 处理完成
    # 或 "RUN_ERROR"         # 错误情况
]

这种设计的好处:

  • 实时反馈:用户立即知道AI开始处理

  • 渐进式显示:逐词输出提升用户体验

  • 错误处理:统一的错误反馈机制

双协议支持:WebSocket + HTTP SSE

python

# WebSocket - 实时双向通信
@app.websocket("/ws/chat/{session_id}")

# HTTP SSE - 兼容性更好的单向流式通信  
@app.post("/api/chat/{session_id}/stream")

第四章:前端开发 - 打造现代化聊天界面

4.1 完整前端HTML代码

由于前端代码较长,这里提供关键部分解析,完整代码已在前文提供。

4.2 前端核心技术解析

AG-UI事件处理机制

javascript

// AG-UI事件处理器
function processAGUIEvent(event) {
    switch (event.type) {
        case 'RUN_STARTED':
            // 显示"思考中"状态
            showTypingIndicator();
            currentRunId = event.run_id;
            break;
            
        case 'TEXT_MESSAGE_CONTENT':
            // 流式更新消息内容
            updateMessageContent(event.delta);
            break;
            
        case 'RUN_FINISHED':
            // 隐藏"思考中"状态
            hideTypingIndicator();
            break;
            
        case 'RUN_ERROR':
            // 显示错误信息
            showError(event.message);
            break;
    }
}
流式响应处理

javascript

// 处理服务器流式响应
async function processStreamResponse(response) {
    const reader = response.body.getReader();
    const decoder = new TextDecoder();
    
    while (true) {
        const { done, value } = await reader.read();
        if (done) break;
        
        const chunk = decoder.decode(value);
        const events = parseAGUIEvents(chunk);
        
        events.forEach(event => {
            processAGUIEvent(event);
            logEvent(event);  // 调试用事件日志
        });
    }
}

第五章:完整部署指南

5.1 一键启动脚本

创建 start.sh(Linux/macOS)或 start.bat(Windows):

start.sh:

bash

#!/bin/bash
echo "🤖 启动AG-UI聊天应用..."

# 检查Ollama服务
echo "检查Ollama服务..."
if ! curl -s http://localhost:11434/api/tags > /dev/null; then
    echo "启动Ollama服务..."
    ollama serve &
    sleep 5
fi

# 检查Python环境
echo "检查Python环境..."
if ! command -v python &> /dev/null; then
    echo "错误:未找到Python,请先安装Python 3.8+"
    exit 1
fi

# 安装依赖
echo "安装Python依赖..."
pip install -r requirements.txt

# 启动FastAPI服务
echo "启动FastAPI服务器..."
python main.py

start.bat:

batch

@echo off
echo 🤖 启动AG-UI聊天应用...

echo 检查Ollama服务...
curl -s http://localhost:11434/api/tags > nul
if errorlevel 1 (
    echo 启动Ollama服务...
    start ollama serve
    timeout /t 5
)

echo 检查Python环境...
python --version > nul 2>&1
if errorlevel 1 (
    echo 错误:未找到Python,请先安装Python 3.8+
    pause
    exit /b 1
)

echo 安装Python依赖...
pip install -r requirements.txt

echo 启动FastAPI服务器...
python main.py

5.2 依赖文件

创建 requirements.txt

text

fastapi==0.104.1
uvicorn==0.24.0
websockets==12.0
requests==2.31.0
ag-ui-protocol>=0.1.0
python-multipart==0.0.6

第六章:测试与调试

6.1 验证部署

  1. 检查Ollama服务

    bash

    curl http://localhost:11434/api/tags
  2. 测试后端API

    bash

    curl http://localhost:8000/health
  3. 访问前端界面
    打开浏览器访问 http://localhost:8000

6.2 常见问题解决

问题1:Ollama连接失败

bash

# 解决方案:
ollama serve
# 或重新安装Ollama

问题2:端口被占用

bash

# 解决方案:更改端口
uvicorn main:app --port 8001

问题3:模型下载失败

bash

# 解决方案:使用国内镜像
OLLAMA_HOST=0.0.0.0 ollama pull gemma2:2b

第七章:进阶功能扩展

7.1 添加更多AI模型支持

python

def call_multiple_ai_services(message: str, provider: str = "ollama") -> List[str]:
    """支持多个AI服务提供商"""
    if provider == "ollama":
        return call_ollama_service(message)
    elif provider == "openai":
        return call_openai_service(message)
    elif provider == "azure":
        return call_azure_openai_service(message)
    else:
        return simulate_ai_thinking(message)

7.2 添加对话记忆

python

def get_conversation_context(session_id: str, max_messages: int = 10) -> str:
    """获取对话上下文"""
    if session_id not in chat_sessions:
        return ""
    
    messages = chat_sessions[session_id][-max_messages:]
    context = "\n".join([f"{msg['role']}: {msg['message']}" for msg in messages])
    return context

总结

通过本文的完整指南,你已经学会了:

✅ 本地部署Ollama AI服务
✅ 理解AG-UI协议的核心概念
✅ 构建完整的FastAPI后端
✅ 开发现代化聊天前端
✅ 实现实时流式聊天体验
✅ 掌握项目部署和调试技巧

AG-UI协议为AI应用开发带来了标准化和高效性,让你能够专注于业务逻辑而不是底层通信细节。现在,你可以基于这个基础项目,扩展更多有趣的功能,如文件上传、多模态对话、语音交互等!


下一步行动:

  1. 🚀 立即运行代码体验效果

  2. 🔧 根据自己的需求定制功能

  3. 🌟 分享你的改进和创意

祝你在AG-UI的世界里构建出精彩的AI应用!

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐