使用 Pydantic AI 的 VercelAIAdapter 对接 Vercel AI Data Stream Protocol 构建流式聊天后端
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
导读
本指南讲解 Pydantic AI 官方提供的 Vercel AI 集成能力:通过VercelAIAdapter将前端(AI SDK UI 的useChat等 hooks)发来的 Vercel AI 请求输入转换为Agent.run_stream_events()的参数,运行 Agent 后把 Pydantic AI 事件流实时转换为 Vercel AI 的 SSE 事件流返回前端。读完本文,你将掌握 Starlette/FastAPI 下的一行式接入、无 Starlette 框架下的手动编排、取消与工具审批、自定义事件数据下发、客户端工具文件回传、消息元数据往返,以及系统提示词与信任模型的完整配置。
协议对接的总体架构
Vercel AI Data Stream Protocol 是 AI SDK UI 生态(useChat、useAssistant等 hooks 与 AI Elements 组件)默认使用的流式协议。Pydantic AI 在pydantic_ai.ui.vercel_ai模块中实现了该协议的双向转换,核心是两个类:
VercelAIAdapter:负责把前端请求体转换为Agent.run_stream_events()的输入、运行 Agent、再把 Pydantic AI 事件转换为 Vercel AI 事件。它继承自抽象的UIAdapter基类(与 AG-UI 协议 的AGUIAdapter平级,二者共用同一套适配器骨架)。VercelAIEventStream:负责把 Pydantic AI 原生事件流转换为 Vercel AI 的 chunk 序列并编码为 SSE 字符串。常规请求路径下你无需直接使用它,只有 Agent 事件不是经由"服务于前端的那个请求"到达你时(例如 durable execution 工作流、消息队列或 WebSocket 扇出场景),才需要单独实例化它来转换事件,参见 UI Event Streams 文档 中的encode_events示例。
从源码看,VercelAIAdapter是一个 dataclass,声明了sdk_version(默认5)、server_message_id等字段,并复用了UIAdapter提供的dispatch_request、from_request、run_stream、run_stream_native、encode_stream、streaming_response、transform_stream等完整方法族。整个请求处理链条为:build_run_input(请求体)→VercelAIAdapter(agent, run_input, accept)→run_stream(...)→encode_stream(...)→StreamingResponse。
与 Starlette / FastAPI 一行式接入
如果后端基于 Starlette 系框架(FastAPI、Litestar 等),VercelAIAdapter.dispatch_request()类方法可以直接在端点函数中消费请求并返回流式响应,这是官方推荐的最简路径:
from fastapi import FastAPI from starlette.requests import Request from starlette.responses import Response from pydantic_ai import Agent from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() @app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request(request, agent=agent)除了request与agent,dispatch_request还接受与Agent.run_stream_events()相同的可选参数(message_history、deferred_tool_results、conversation_id、run_id、model、instructions、deps、output_type、model_settings、usage_limits、usage、metadata、infer_name、toolsets、capabilities等),以及两个回调:
on_complete:Agent 运行成功时触发,接收AgentRunResult,可额外产出 Vercel AI 事件;on_cancel:Agent 因一等取消(first-party cancellation)结束时触发,接收RunCancelled,同样可额外产出事件。
调用链在 UIAdapter.dispatch_request 中实现:先from_request()解析请求,若请求体校验失败则返回 422(UNPROCESSABLE_ENTITY)JSON 错误,随后streaming_response(adapter.run_stream(...))生成流式响应。注意请求的Accept头会被用作流式响应的media_type。
无 Starlette 框架时的手动编排(Advanced Usage)
对于 Django、Flask 等非 Starlette 框架,或需要细粒度控制输入输出的场景,可以直接实例化VercelAIAdapter并链式调用其方法,效果等价于dispatch_request:
- 解析请求体:
VercelAIAdapter.build_run_input(await request.body())返回 Vercel AI 的RequestData对象。其实现是request_data_ta.validate_json(body)(见 _adapter.py),即用 Pydantic TypeAdapter 严格校验 JSON。如果是 Starlette/FastAPI 环境,也可直接用VercelAIAdapter.from_request()一步构建适配器实例。 - 运行 Agent:
adapter.run_stream(...)运行 Agent 并返回 Vercel AI 事件流,支持与Agent.run_stream_events()相同的参数及on_complete/on_cancel。它内部先走run_stream_native()(返回 Pydantic AI 原生事件),再经transform_stream()转换;你也可以自己先拿run_stream_native()的原生事件流,再调用transform_stream()手动转换。 - 编码 SSE:
adapter.encode_stream(event_stream)把 Vercel AI 事件流编码为 SSE 字符串;或者直接用adapter.streaming_response(...)生成 Starlette/FastAPI 的流式响应对象。
完整的可运行示例(含输入校验、422 返回、取消 token 注册表清理)参见官方文档 run_stream.py 中的对应示例与下方"取消"一节。该示例使用 FastAPI,但可改造适配任意 Web 框架。
取消机制:一等取消与外部取消
当一个运行以一等取消结束(来自工具内ctx.cancel()、AgentRun.cancel()、或你服务器端在取消端点接线的CancellationToken),适配器会向流中发射一个 Vercelabortchunk。useChat会保留已产生的部分消息,并在onFinish中上报isAbort,而不是进入错误状态。此时on_cancel回调被触发,可在其中持久化可恢复的消息历史(RunCancelled.all_messages()返回可继续传递的消息列表)。
重要区别:客户端调用stop()会直接中断浏览器请求,服务器侧看到的是一次"断开"(disconnect),属于外部取消——运行以asyncio.CancelledError方式被撕裂,不会发射abortchunk,on_cancel也不会触发(客户端反正已经断开了)。若希望在用户点"停止"时也能拿到abortchunk 并执行on_cancel,应当保持流连接,改用一等取消:给运行传入CancellationToken,并暴露一个独立端点(如POST /chat/{id}/cancel)调用token.cancel()。完整实现如下:
import json from collections.abc import AsyncIterator from http import HTTPStatus from fastapi import FastAPI from fastapi.requests import Request from fastapi.responses import Response, StreamingResponse from pydantic import ValidationError from pydantic_ai import Agent, CancellationToken, RunCancelled from pydantic_ai.ui import SSE_CONTENT_TYPE from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() cancellation_tokens: dict[str, CancellationToken] = {} async def on_cancel(cancelled: RunCancelled) -> None: messages = cancelled.all_messages() # (1)! print(f'cancelled after {len(messages)} messages') @app.post('/chat/{chat_id}') async def chat(chat_id: str, request: Request) -> Response: accept = request.headers.get('accept', SSE_CONTENT_TYPE) try: run_input = VercelAIAdapter.build_run_input(await request.body()) except ValidationError as e: return Response( content=json.dumps(e.json()), media_type='application/json', status_code=HTTPStatus.UNPROCESSABLE_ENTITY, ) adapter = VercelAIAdapter(agent=agent, run_input=run_input, accept=accept) cancellation_token = CancellationToken() cancellation_tokens[chat_id] = cancellation_token event_stream = adapter.run_stream( cancellation_token=cancellation_token, on_cancel=on_cancel ) async def encode_stream() -> AsyncIterator[str]: try: async for event in adapter.encode_stream(event_stream): yield event finally: if cancellation_tokens.get(chat_id) is cancellation_token: cancellation_tokens.pop(chat_id, None) return StreamingResponse(encode_stream(), media_type=accept) @app.post('/chat/{chat_id}/cancel', status_code=HTTPStatus.NO_CONTENT) async def cancel_chat(chat_id: str) -> None: if token := cancellation_tokens.get(chat_id): token.cancel()- 这是需要持久化的可恢复历史——在后续运行中作为
message_history传入即可续接对话。
注意:上述内存 token 注册表要求单进程部署或粘性路由。多 worker 部署时,需借助消息代理等共享协调机制,把取消请求路由到持有该运行的 worker。
从源码看,VercelAIEventStream.on_cancelled()发射的AbortChunk内容为reason='The agent run was cancelled.'(见 _event_stream.py);而on_error()会设置finish_reason='error'并发射ErrorChunk(_event_stream.py)。
向客户端下发数据:Data Chunks
运行过程中(例如长耗时工具的执行进度)想向客户端推送自定义数据,可以在工具内通过ctx.emit()发射一个CustomEvent:
from dataclasses import dataclass from pydantic_ai import Agent, CustomEvent, RunContext agent = Agent('openai:gpt-5.2') @dataclass(kw_only=True) class FileUploadProgressEvent(CustomEvent): done: int total: int @agent.tool async def upload_files(ctx: RunContext, total: int) -> str: for done in range(1, total + 1): # 完成一个单位的工作,然后告诉前端进度 await ctx.emit(FileUploadProgressEvent(done=done, total=total)) return f'Uploaded {total} files'每个事件会以DataChunk形式到达客户端,type为data-{name},data为to_payload()的结果——上例即type='data-file_upload_progress'、data={'done': 1, 'total': 3}。chunk 在事件发射的当下、工具仍在运行时即到达,天然支持进度条。
形状一致性:无论事件是否从工具调用内部发射,data形状都一致,因此前端针对某一形状编写的代码,不会因为同一事件类日后改在别处发射而失效。如需自定义形状(比如按前端期望命名字段、把工具归属信息放到线上),覆盖to_payload()即可:
from dataclasses import dataclass from typing import Any from pydantic_ai import CustomEvent @dataclass(kw_only=True) class FileUploadPhaseEvent(CustomEvent): done: int total: int def to_payload(self) -> dict[str, Any]: return { 'completed': self.done, 'total': self.total, 'toolCallId': self.tool_call_id, }特殊规则:若to_payload()返回一个数据承载型 chunk(见下文),则该 chunk 原样透传;事件类声明为ui=False时永不下发(仅服务端消费的事件留在服务端);进程从未 import 过的事件类也不会被转发,因为其 opt-out 标记是挂在类上而非走线上。
工具返回时携带 chunk:工具可通过返回带metadata(单个或列表)的ToolReturn对象,把 Vercel AI data stream chunk 附加到工具结果上。支持四种类型:DataChunk、SourceUrlChunk、SourceDocumentChunk、FileChunk:
from pydantic_ai import Agent, ToolReturn from pydantic_ai.ui.vercel_ai.response_types import DataChunk, SourceUrlChunk agent = Agent('openai:gpt-5.2') @agent.tool_plain async def search_docs(query: str) -> ToolReturn: return ToolReturn( return_value=f'Found 2 results for "{query}"', metadata=[ SourceUrlChunk( source_id='doc-1', url='https://example.com/docs/intro', title='Introduction', ), DataChunk( type='data-search-results', data={'query': query, 'count': 2}, ), ], )与ctx.emit发射的事件不同,这些 chunk 属于消息的一部分,能随消息历史往返而幸存——这正是前端需要重建的数据(如答案背后的来源 URL)所期望的;代价是它们在工具返回时而非运行过程中下发。源码中iter_metadata_chunks()(见 _utils.py)只会转发上述四种>@app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request(request, agent=agent, sdk_version=6)
当sdk_version=6时适配器会:
- 在调用
requires_approval=True的工具时发射tool-approval-requestchunk; - 自动从后续请求中提取审批响应;
- 为被拒绝的工具发射
tool-output-deniedchunk。
前端方面,AI SDK UI 的useChathook 处理审批流程;可使用 AI Elements 的Confirmation组件做现成审批 UI,或用 hook 的addToolApprovalResponse自行构建。
审批响应的严格性:审批响应按协议经useChat的addToolApprovalResponse与参考 Next.js 后端往返,设计上被信任。但审批决定本身必须是真正的 JSON 布尔值——ToolApprovalResponded.approved字段是严格布尔(StrictBool),任何替代值(1、"true"、0、"false")都会导致请求校验失败,而不会被强制转换成一个决定(见 request_types.py 的注释,防止{'approved': 1}被宽松模式强转为通过)。若需要把审批决定绑定到服务端状态而非请求本身,可拦截DeferredToolRequests,在服务端持久化审批 ID,并在续接时显式传入deferred_tool_results。
补充说明:源码中sdk_version类型为Literal[5, 6, 7],默认5以保证向后兼容;7的线协议与6完全相同(v7 的>from fastapi import FastAPI from starlette.requests import Request from starlette.responses import Response from pydantic_ai import Agent from pydantic_ai.ui.vercel_ai import VercelAIAdapter agent = Agent('openai:gpt-5.2') app = FastAPI() @app.post('/chat') async def chat(request: Request) -> Response: return await VercelAIAdapter.dispatch_request( request, agent=agent, manage_system_prompt='client' )
协议细节的源码级补充
- Finish reason 映射:Pydantic AI 的结束原因经
_FINISH_REASON_MAP(_event_stream.py)映射为 Vercel 格式:stop→stop、length→length、content_filter→content-filter、tool_call→tool-calls、error→error,未知值统一为other。 - SSE 编码与响应头:每个 chunk 编码为
data: {json}\n\n(encode_event),且流式响应带x-vercel-ai-ui-message-stream: v1响应头(VERCEL_AI_DSP_HEADERS),这是 AI SDK UI 识别 data stream 协议的标准信号。 - 消息 ID 生成:
_generate_message_id按优先级生成确定性消息 ID:有provider_response_id用{provider_response_id}-{index},有run_id用{run_id}-{index},否则用uuid5(timestamp-kind-role-index)(_adapter.py)。 conversation_id关联:适配器把请求体顶层的id(chat ID)作为conversation_id,用于跨多次运行关联 OpenTelemetry span 与消息历史。- 往返的已知损耗:
dump_messages → load_messages对工具结果并非完全无损——RetryPromptPart重新加载后会变成outcome='failed'的ToolReturnPart(协议没有独立 retry 概念);无法解析为 JSON 对象的ToolCallPart.args会被重写为{'INVALID_JSON': '<raw args>'}。需要 retry 语义跨往返存活时,应在进程内维护对话而非经由 Vercel AI 线格式持久化。
验证与测试
仓库的测试套件(tests/test_ui.py)覆盖了适配器的请求解析、事件流编码与消息往返;tests/test_vercel_ai.py与tests/cassettes/test_vercel_ai/中的 VCR 磁带用于验证真实模型调用下的流式行为。测试同时覆盖了sdk_version=5/6两种线格式下的 chunk 输出差异,可作为理解协议行为的补充参考。
总结
VercelAIAdapter把 Pydantic AI 的 Agent 运行无缝接入 Vercel AI Data Stream Protocol 生态:dispatch_request一行式接入 FastAPI,手动编排适配任意框架,取消、自定义数据、工具审批、消息元数据与压缩均开箱即用。安全方面,系统提示词归属、文件 URL scheme 白名单与上传文件开关共同构成默认信任模型;在对外暴露该端点时,请务必把它当作内部后端服务,放在你自己已认证的路由处理器之内。
【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考