作者结合自研 chanx 与 Pydantic AI,实现类型安全的 WebSocket 流式 Agent 交互,填补了官方 SSE 方案的 UI 层缺口。
我在多年间一直在构建实时产品:群聊、语音管线、流式 AI 助手。有两种痛点始终如影随形。
第一种是 WebSocket。每个项目的 receive_json 处理函数都长成同样的无类型模样——一串 if/else、手写的验证逻辑,缺少前端可以依赖的契约。消息是字符串类型的,bug 只在运行时才暴露,新队友上手意味着"通读消费者代码,从头到尾"。这种痛促使我构建了 chanx:类型化消息、自动路由、AsyncAPI 文档(完整故事在我的介绍帖里)。
第二种是智能体框架。我最早在 LangChain 上跑 LLM 功能,花在调试它的抽象层上的时间比写自己的代码还多。第一次用 Pydantic AI 就被圈粉了:类型化工具、类型化输出、类型化依赖注入。它将 chanx 应用于 WebSocket 的同样理念应用到了智能体,而且两者都基于 Pydantic,几乎不需要粘合剂就能衔接。
但还有一个缺口。在 UI 层,pydantic-ai 的官方方案是 SSE 形态的:AG-UI 和 Vercel AI 协议的适配器,以及一个内置的 Agent.to_web() 聊天页面,文档明确标注为开发工具。真正的智能体产品需要的不仅是一条单向流。对话是双向的:用户发送提示词,智能体流式返回文本和工具活动,而且(有趣的部分)智能体有时需要在运行破坏性工具之前停下来请求许可。这种往返(暂停、询问、恢复)完美契合 WebSocket,而目前 pydantic-ai 文档中根本没有 WebSocket 示例。
所以我构建了一个,用我构建 REST API 一样的方式:契约优先。服务器将每条消息定义为一个 Pydantic 模型,chanx 将这些模型转化为路由处理器加 AsyncAPI 3.0 schema,React 客户端从该 schema 生成 TypeScript 类型。后端和前端不会产生偏差,队友从契约集成而不是读你的 Python,整个流程(流式、工具调用、审批)无需调用 LLM 即可测试。
配套仓库包含所有可运行的代码:huynguyengl99/pydantic-ai-ws-agent。下方的代码片段为便于阅读做了精简,截至 pydantic-ai 2.9 和 chanx 2.8 均有效;仓库才是权威来源。
演示智能体("Tasklet")在 SQLite 中管理任务列表。安全工具自由运行;破坏性工具则暂停等待审批。

pydantic-ai Agent ──run_stream events──▶ AgentHarness (background task)
│ │ broadcast_event
typed tools ▼
requires_approval=True group: conversation.{id} ◀── HTTP POST /notify
│ channel layer (in-memory / Redis)
┌────────────────┼────────────────┐
▼ ▼ ▼
chanx consumer chanx consumer chanx consumer AsyncAPI 3.0 schema
(tab 1) (tab 2) (other device) ──▶ generated TS types
四个部分:一个普通的 pydantic-ai 智能体(类型化工具、类型化依赖项、危险操作设置 requires_approval、结构化输出);一个 Harness 来运行它并将每个事件转换为类型化消息;一个通道层将这些事件传送到每个对话的组而不是单个 socket;以及一个 chanx 消费者,验证、路由并记录 React 客户端从中生成类型的协议。
智能体:知道自己是危险品的工具
这里没有任何 WebSocket 感知逻辑。这是一个标准的 pydantic-ai 智能体,有两个刻意为之的选择:一个结构化输出类型,以及携带对话 ID 的类型化依赖项,使每个工具都在该对话的任务上操作。
class TextAnswer(BaseModel):
"""Final answer to the user, with suggested follow-up prompts."""
content: str
follow_ups: list[str] = Field(default_factory=list)
@dataclass
class AgentDeps:
conversation_id: str
# Every run ends as either a final answer or a request for tool approval.
AgentOutput = TextAnswer | DeferredToolRequests
def build_agent(model: Model | str | None = None) -> Agent[AgentDeps, AgentOutput]:
agent: Agent[AgentDeps, AgentOutput] = Agent(
model or settings.resolved_model,
deps_type=AgentDeps,
output_type=[TextAnswer, DeferredToolRequests],
instructions=INSTRUCTIONS,
)
@agent.tool
async def add_task(ctx: RunContext[AgentDeps], title: str) -> db.TaskItem:
"""Add a new task with the given title."""
return await db.add_task(ctx.deps.conversation_id, title)
@agent.tool(requires_approval=True)
async def delete_task(ctx: RunContext[AgentDeps], task_id: int) -> str:
"""Permanently delete the task with the given id. Requires user approval."""
await db.delete_task(ctx.deps.conversation_id, task_id)
return f"Task {task_id} deleted"
return agent
(仓库中有六个工具;clear_all_tasks 是另一个需要审批的工具。)
requires_approval=True 是 pydantic-ai 的延迟工具特性。当模型调用 delete_task 时,工具不会执行,而是提前结束运行,result.output 是一个携带待处理调用的 DeferredToolRequests,这正是 output_type 是联合类型的原因:每次运行要么以最终答案结束,要么以权限请求结束,类型系统迫使你处理这两种情况。
TextAnswer 而非普通 str 是结构化输出的另一半。除了答案文本,模型还会用三个建议的后续提示填充 follow_ups,UI 将它们渲染为可点击的标签 chip。最终答案不是附加了约定的散文,而是一个经过验证的对象,chip 标签随同其 schema 一起传递。
build_agent() 作为工厂函数而非模块级单例,在后续测试中很重要:测试用假模型构建相同的智能体。
每条 WebSocket 消息都是一个带有 Literal action 字段的 Pydantic 模型。chanx 用该字段作为验证、路由和文档的判别器:
class ChatMessage(BaseMessage):
"""User sends a prompt to the agent."""
action: Literal["chat"] = "chat"
payload: ChatPayload
完整协议很小。客户端 → 服务器:chat、tool_decision。服务器 → 客户端:user_message、stream_start、text_delta、tool_call、tool_result、approval_request、stream_end、suggestions、history、tasks_updated、notification、agent_error。
消费者将各部分绑定在一起:处理器声明它们接受什么以及可能回复什么,而这个声明就是文档:
@channel(name="agent", description="Pydantic AI task assistant over a typed WebSocket")
class AgentConsumer(BaseConsumer):
passthrough_events: ClassVar[list[type[BaseMessage]]] = BROADCAST_MESSAGES
@ws_handler(
summary="Send a prompt to the agent",
output_type=AgentServerMessage, # the union of everything we may send back
)
async def handle_chat(self, message: ChatMessage) -> None:
self._spawn_run(self.harness.run(message.payload.text))
@ws_handler(summary="Approve or deny pending tool calls", output_type=AgentServerMessage)
async def handle_tool_decision(self, message: ToolDecisionMessage) -> None:
self._spawn_run(self.harness.resolve(message.payload.decisions))
没有 if/else 路由,没有手动验证,/asyncapi.json 现在提供一份 AsyncAPI 3.0 规范,描述上述每条消息:这是 WebSocket 版的 drf-spectacular 为 REST 提供的功能(我在契约优先的帖子里写过这个工作流)。
流式结构化输出:将快照转为增量
流式传输纯文本很简单。流式传输结构化答案才是有趣的部分:最终输出是一个 TextAnswer 对象,所以没有原始文本流可以转发。pydantic-ai 的解决方案是 run_stream() 加上两个钩子。
工具事件通过 event_stream_handler 到达,Harness 将其直接映射到消息契约:
async def forward_tool_events(
_ctx: RunContext[AgentDeps], events: AsyncIterable[AgentStreamEvent]
) -> None:
async for event in events:
match event:
case FunctionToolCallEvent(part=part):
await self.send(ToolCallMessage(payload=ToolCallPayload(
tool_call_id=part.tool_call_id,
tool_name=part.tool_name,
args=part.args_as_dict(),
)))
case FunctionToolResultEvent(part=part):
await self.send(ToolResultMessage(payload=ToolResultPayload(
tool_call_id=part.tool_call_id,
content=to_jsonable_python(part.content),
)))
答案本身以部分验证的快照形式流式传输:stream_output() 在模型的 JSON 每次到达更多内容时都 yield 一个不断增长的 TextAnswer。对比连续的内容值可以恢复要发送的文本增量:
sent_text = ""
async with self.agent.run_stream(
user_prompt,
message_history=history,
deferred_tool_results=deferred_tool_results,
deps=AgentDeps(conversation_id=self.conversation_id),
event_stream_handler=forward_tool_events,
) as stream:
async for partial in stream.stream_output(debounce_by=None):
match partial:
case TextAnswer(content=content) if content and content != sent_text:
delta = (
content.removeprefix(sent_text)
if content.startswith(sent_text)
else content
)
await self.send(TextDeltaMessage(payload=TextDeltaPayload(delta=delta)))
sent_text = content
output = partial
await db.save_history(self.conversation_id, stream.all_messages_json().decode())
客户端会看到用于散文的 text_delta,以及用于工具活动的 tool_call / tool_result 对(以实时卡片形式在 UI 中渲染),而循环本身完全不感知存在哪些工具。
循环结束后,联合输出类型充分发挥了其价值:
match output:
case DeferredToolRequests() as requests:
# run_stream 在找到最终结果(暂停的调用)后停止转发事件,
# 所以在这里宣布它们:相同的 id,相同的传输顺序。
for part in requests.approvals:
await self.send(ToolCallMessage(payload=ToolCallPayload(...)))
PENDING_APPROVALS[self.conversation_id] = requests
await self.send(approval_request_message(requests))
case TextAnswer(content=content, follow_ups=follow_ups):
await self.send(StreamEndMessage(payload=StreamEndPayload(text=content, usage=...)))
if follow_ups:
await db.save_suggestions(self.conversation_id, follow_ups)
await self.send(SuggestionsMessage(payload=SuggestionsPayload(follow_ups=follow_ups)))
有一个值得指出的细节:延迟的工具调用永远不会到达 event_stream_handler(延迟意味着"未执行",只有正在执行的工具才会发送事件),所以 harness 会在 approval_request 之前自行宣布暂停的调用。客户端以相同的方式渲染工具卡片;它无法区分是哪条路径产生的。而在正常流程中,建议以自己的消息形式在 stream_end 之后到达,所以回答稳定后建议标签会弹入。
广播,而非发送:后台运行与分组
我的第一个版本在聊天处理器内 await agent 运行,并将消息直接写入 socket。这能工作,但有三个隐蔽的问题:处理器在整个运行期间阻塞、页面刷新 mid-run 会中断流、第二个标签页看不到任何内容。
修复方案是 chanx(以及更早的 Django Channels)所围绕的模式:广播到分组,而非连接。每个消费者在连接时加入一个按会话 ID 分组的组,而运行作为后台任务执行,向该组广播事件:
async def post_authentication(self) -> None:
group = f"conversation.{self.conversation_id}"
await self.channel_layer.group_add(group, self.channel_name)
self.harness = AgentHarness(agent, self.conversation_id, self._broadcast)
async def _broadcast(self, message: BaseMessage) -> None:
await self.broadcast_event(message, groups=f"conversation.{self.conversation_id}")
def _spawn_run(self, coro) -> None:
# 该任务的生命周期长于此连接:在 mid-run 时刷新的客户端
# 仍然能通过会话组接收事件。
RUNNING[self.conversation_id] = asyncio.create_task(guarded(coro))
在接收端无需编写任何代码。chanx 的 passthrough_events 会生成事件处理器:通过组收到的列表中的任何消息类型都会逐字转发到 WebSocket 客户端,并与所有其他内容一起记录在 AsyncAPI 中。
具体来说,这带来了什么(全部有测试覆盖):
聊天处理器立即返回;按会话 ID 的 RUNNING 字典用类型化的 agent_error 拒绝并发运行
刷新安全流:运行在没有任何 socket 附着的情况下继续;重连会重放历史并获取实时事件
多标签页开箱即用:每个标签页渲染相同的 delta、工具卡片和审批请求。甚至用户自己的提示也通过这种方式传播——在 stream_start 之前作为 user_message 回显,所以发送标签页从与其他人相同的广播中渲染气泡,而不是在本地追加
任意位置审批:待审批项按会话而非连接进行键控;在一台设备上请求审批,在另一台上批准
超越 Socket:HTTP、Worker、其他框架
broadcast_event 是一个类方法,所以任何东西都可以向会话中推送类型化消息。示例公开了 POST /conversations/{id}/notify,而 curl 可以向每个连接的标签页发送 toast。agent 从内部使用同一扇门:一个 schedule_reminder 工具从 ctx.deps 读取会话 ID、调度作业并立即返回;当作业在 stream_end 很久之后触发时,它会向该组广播通知。
这扇逃生舱可以向两个方向扩展。组广播不是唯一的投递模式:chanx 可以寻址单个连接,所以处理器可以将繁重工作交给真正的任务队列(taskiq、Celery、ARQ),使用 self.channel_name 作为类型化回复地址:
# consumer: 附带回复地址将作业入队
await run_agent_job.kiq(prompt, reply_to=self.channel_name)
# worker 进程: 将类型化事件发送回那个确切的连接
await AgentConsumer.send_event(JobProgressEvent(payload=...), reply_to)
# consumer: 类型化、路由化,与其他内容一样记录在案
@event_handler
async def handle_job_progress(self, event: JobProgressEvent) -> ProgressMessage:
return ProgressMessage(payload=event.payload)
consumer 保持轻量级协议适配器,而 LLM 工作发生在为它构建的进程中。由于 channel 层只是 Redis,而 chanx 在 Django Channels 和 FastAPI 上使用相同的 API,socket 和 agent 甚至不需要共享服务。我在生产环境中运行的一种拆分:Django 拥有面向用户的 WebSocket(session、ORM auth、一个几乎只是 passthrough_events 的 consumer),而独立的 FastAPI 服务运行模型并广播到相同的会话组。Django 从不 await LLM;agent 服务从不触碰 session。它们共享的只有契约,而契约部分由 AsyncAPI 记录。
演示运行在零依赖的内存层上;设置 REDIS_URL 相同的代码即可跨进程和实例扩展。
人在环中:审批往返
这是 WebSocket 生来就要处理的场景。完整序列:
delete_task → 运行以 DeferredToolRequests 结束,harness 宣布暂停的调用为 tool_call 消息,所以 UI 立即显示工具卡片(状态:等待审批)approval_request,持久化消息历史,然后直接停止。无轮询,无挂起请求tool_decisionDeferredToolResults 并用保存的历史恢复运行:async def resolve(self, decisions: list[ToolDecision]) -> None:
approvals = {
d.tool_call_id: (
ToolApproved(override_args=d.override_args)
if d.approved
else ToolDenied(message=d.reason or "Denied by the user")
)
for d in decisions
}
del PENDING_APPROVALS[self.conversation_id]
await self.run(None, DeferredToolResults(approvals=approvals))
恢复的运行只是另一个 run_stream() 调用:已批准的工具执行(它们的事件通过相同的处理器流式传输),被拒绝的工具向模型返回拒绝消息,模型在任一情况下都流式传输最终响应。ToolApproved(override_args=...) 是一个很棒的细节:用户可以在批准时编辑参数,这才是真正的审查 UI 应该有的工作方式。
审批在两个方向上都能抗重连。待处理请求按会话而非连接键控,consumer 向每个新连接重新发送开放的 approval_request,所以 mid-approval 时的刷新会得到一个可工作的审批卡片而非死路一条。端到端测试实际上是从第二个浏览器标签页批准的。一个 UI 细节来自多标签页:当运行恢复时,其他仍显示开放审批卡片的标签页会在下一个 stream_end 时将其标记为"已从其他会话解决"。
陷阱:告诉模型 harness 会处理它
我的第一条指令是"破坏性操作需要用户审批"。模型把它当成自己的事了。没有调用 delete_task,而是在聊天中问:"这是永久性的,我得到您的批准了吗?(是/否)"。技术上顺从,完全绕过了类型化的审批流程。
修复方法是明确分工:
当用户请求一个破坏性操作时,直接调用工具——不要在聊天中请求确认。系统会自动暂停破坏性工具调用并在 UI 中请求用户批准。
如果构建了审批流程,这个提示词 bug 会找上你。模型需要知道安全是基础设施,不是它的工作。
AsyncAPI → TypeScript:客户端自动生成
chanx 在 /asyncapi.json 提供契约;上述每个 payload 都作为普通 JSON Schema 出现在 components.schemas 中。一个约 100 行的脚本获取 spec,从 AsyncAPI 操作中收集客户端绑定和服务器绑定的消息(输入 vs 回复),然后将它们喂给 json-schema-to-typescript:
// generated - src/generated/messages.ts
export interface ToolCallMessage {
action: "tool_call";
payload: ToolCallPayload;
}
export type ClientMessage = ChatMessage | ToolDecisionMessage;
export type ServerMessage = /* 所有十二个服务器消息的联合类型 */;
一个值得了解的 codegen 细节:Pydantic 在 JSON Schema 中将 action 标记为可选的(它有默认值),但它是鉴别器(discriminator),所以脚本强制将其放入 TypeScript 的 required 中以便进行类型收窄。有了这个,整个客户端状态层就是一个穷举式 switch:
function applyServer(state: ChatState, msg: ServerMessage): ChatState {
switch (msg.action) {
case "text_delta": // 追加到流式气泡
case "tool_call": // 添加工具卡片
case "approval_request": // 渲染批准/拒绝按钮
case "suggestions": // 显示后续建议芯片
case "stream_end": // 最终确定 + 统计 token
// TypeScript 会在任何服务器消息未处理时报错
}
}
在 Pydantic payload 中重命名字段,运行 pnpm generate,tsc 会列出所有刚刚变得不正确的客户端行。schema-first REST 的同一循环:传输层变了,纪律没变。
持久化:历史记录作为类型化产物
pydantic-ai 提供了 ModelMessagesTypeAdapter,一个用于完整消息历史的 TypeAdapter,因此会话状态只是另一个经过验证的 blob:stream.all_messages_json() 出去,ModelMessagesTypeAdapter.validate_json(raw) 回来。两个函数,没有自定义序列化。
连接时,消费者重放会话作为历史消息加上当前任务列表(tasks_updated)和最新的后续建议芯片。转录构建器遍历类型化历史并重建交错顺序:用户文本、助手文本和工具卡片按顺序排列,每张卡片带有 done、denied 或 awaiting 状态,因此在刷新后被拒绝的删除仍会渲染为已拒绝。由于所有内容都按会话存储,带有 /conversations 的侧边栏(标题来自第一个提示)和会话删除几乎是免费获得的。多轮上下文、断线重连和批准后恢复都基于这两个函数。
无需 LLM 即可测试完整流程
我最关心的部分:整个 WebSocket 协议(流式、工具执行、批准、拒绝、历史重放)在 pytest 中运行,无需 API key。pydantic-ai 的 FunctionModel 扮演模型;一个 stream 函数脚本精确描述"LLM"的行为。结构化输出的一个波折:最终答案本身就是对 output 工具的调用,所以脚本也会发出它:
async def delete_task_flow(messages, info):
tool_return = next(
(p for p in messages[-1].parts if isinstance(p, ToolReturnPart)), None
)
if tool_return is not None:
# 工具运行后,通过结构化输出工具回答
yield {0: DeltaToolCall(
name=info.output_tools[0].name,
json_args=json.dumps({"content": f"Result: {tool_return.content}"}),
)}
else:
yield {0: DeltaToolCall(name="delete_task", json_args='{"task_id": 1}')}
chanx 的 WebsocketCommunicator 驱动消费者,断言的写法就像协议规范:
await comm.send_message(ChatMessage(payload=ChatPayload(text="Delete task 1")))
replies = await comm.receive_all_messages(stop_action="approval_request")
# 提示回显到群组,暂停的调用被宣布,运行暂停
assert [m.action for m in replies] == [
"user_message", "stream_start", "tool_call", "approval_request",
]
assert (await db.list_tasks("conv-delete")) != [] # 尚未删除任何内容
await comm.send_message(ToolDecisionMessage(...approved=True...))
replies = await comm.receive_all_messages(stop_action="stream_end")
assert await db.list_tasks("conv-delete") == [] # 现在它没了
其余测试覆盖了拒绝路径、两个标签页接收相同广播流、在断开连接后从新连接批准、并发运行守卫、带拒绝和等待工具卡片的转录重放,以及一个契约测试——拉取 /asyncapi.json 并断言每种消息类型都存在,这样前端生成所依赖的 schema 就不会静默丢失任何消息。二十个测试,几秒钟,无需 LLM。
一个值得分享的陷阱:Agent.override(model=...) 使用 contextvars,它不会传播到测试通信器下消费者的任务中。工厂模式绕开了它:测试调用 build_agent(FunctionModel(...)) 并直接打补丁消费者的 agent。
该演示有三个你会在生产环境中重新审视的简化:
横向扩展:仓库提供了一个 docker-compose.yml 用于 Redis,README 逐步讲解了两个实例的演示(聊天在一个端口,在另一个端口观看同一个流)。
认证:消费者信任会话查询参数。chanx 有一个身份验证钩子(post_authentication),你可以在那里验证会话或令牌,然后才加入群组并重放历史。
持久化待批准项:待处理的 DeferredToolRequests 活在进程内字典中;如果批准必须在重启后存活,将它们存储在数据库中。而且如果你曾经接受过来自客户端的消息历史,先用 pydantic-ai 的 sanitize_messages 处理它。
WebSocket 是带批准的智能体的天然传输层:暂停和恢复就是消息。