上一篇把 HTTP、WebSocket、SSE 的方向差讲清楚了:问一次答一次用 HTTP,服务端一直推用 SSE,两头实时互发用 WebSocket。那篇文章里有一句当时完全成立的话——AI 流式回答「用户并不需要在同一条连接上反复给模型发东西」。
Agent 框架往前走了一步,这句话开始不够用了。
用户看着模型往歪了,想立刻改方向;Agent 要删文件、要调支付、要开权限,必须停下来等人点头。这不是另起一段用户任务;但在执行层面,它往往意味着结束当前 generation,并在同一任务上下文中开启新的 run。用户感觉是「同一次任务里插了一句」,后端则是「旧 run 停掉,新 run 接着同一份上下文」。
WebSocket 双向,插话是本职工作。SSE 单向,很多人就默认:要插话就得把当前流掐掉,再开一条。用过就知道,那种感觉很差——字正在往外蹦,页面闪一下,进度条归零,后端那趟 run 也不知道算不算还活着。
其实问题不在 SSE 本身。真正别扭的是:把「用户这句话」和「这条事件流」焊在了同一次 HTTP 请求上。
拆开就顺了。用户输入走普通 HTTP,后端先还一个应用层会话 ID;前端拿这个 ID 去挂 SSE。授权也好,再次引导也好,都是往同一个会话里再 POST 一次。不需要因为用户插话或授权而主动拆流;连接自己掉了,再按事件 ID 续播。
早期 Agent:问一次,流一段,结束
最早一批 Chat UI,交互模型非常干净:
用户提交一句话
POST /chat (响应直接是 text/event-stream)
token 一段段推回来
模型说完,连接关掉
这跟普通补全 API 几乎是同一个形状,只是把整段 JSON 换成了 SSE。前端用 fetch 读流,或者用 EventSource 挂上去,把 data: 拼到对话框里。单向,刚好。
那时候 Agent Loop 也短。Agent Loop 那篇里写过:模型点名,循环执行,再问,停机。用户通常只在开头开口,中间最多看一眼工具调用卡片,并不参与。
所以「一次提问 = 一条 SSE」成了很多项目的默认实现。它简单,本地 demo 也好看。
然后产品开始要两件事。
中途插话和授权,把单向流逼到墙角
第一件是纠偏。模型已经开始写方案了,用户一眼看出方向不对:「别写 React,用 Vue。」「先别改代码,先列风险。」如果等它说完再发下一条,token 白烧,体验也像在跟一个听不见人的人说话。
第二件是授权。Agent 要执行有副作用的工具:覆盖文件、发邮件、调生产接口、读隐私数据。这时候循环必须停住,把「我准备干什么」推给前端,等人点允许或拒绝,再继续。这就是人机协同(HITL)最常见的一种。
WebSocket 处理这两件事更直接:
同一条 WSS
├─ 服务端:token / 工具状态 / 等授权
└─ 客户端:纠偏 / 同意 / 拒绝
双向消息走同一条连接,少一次「再开一条流」的心智负担。但会话是否持续、断线后如何恢复,仍然是应用层状态管理问题。WebSocket 断了,如果 run、审批、事件还只活在内存里,照样要自己恢复——它并不比 SSE 更「自带状态」。
SSE 没有这条回程。规范就是服务端往下推,浏览器不能在同一条 EventSource 上再塞一个 body 上去。于是很多实现会做成这样:
用户点「允许删除」
前端 abort 当前 EventSource
再 POST 一次,带上 approval
后端新开一条 SSE
前端重新订阅、重新拼状态
能跑。但有三处硌手:
- 观感断了。流式输出最怕的就是中间闪一下。重连、重挂、重新对齐光标,用户会觉得系统卡了或者重来了。
- run 的生命周期被绑死在连接上。连接一断,后端如果把这次 Agent 循环一起拆掉,授权到来时那个「等在工具门口」的协程已经没了;如果不断,你又得自己发明一套「断流但 run 还活着」的状态机——而这套状态机,本来就不该依赖连接。
- 前端状态难交接。旧流里已经推了半段文本、两次工具调用、一张权限卡片。新流从哪接着播?重放全部?只播增量?漏一条事件,UI 就会和后端各活各的。
所以别被「SSE 不能回传」吓到。浏览器本来就能同时开很多 HTTP 请求。SSE 只是不能在这一条连接上反向说话,并不禁止你再发一个 POST。别扭的是那个默认假设:这条流是为「这一句用户输入」而生的,输入没了,流就该死。
该拆开的是命令和订阅,不是换协议
把一次 Agent 任务看成两个通道:
| 通道 | 方向 | 干什么 | 用什么 |
|---|---|---|---|
| 命令通道 | 客户端 → 服务端 | 提问、纠偏、授权、取消 | 普通 HTTP POST |
| 事件通道 | 服务端 → 客户端 | token、工具开始/结束、等授权、run 结束 | SSE |
它们共享的不是某次 HTTP 请求,而是一个应用层交互会话。这里的 conversation_id 不是登录 Cookie 里的 Session,也不是一次模型请求。它只是把「这一次用户任务」的命令和事件对上号。
交互变成:
1. POST /conversations
body: { "command_id": "cmd_...", "message": "帮我把登录页改成 Vue" }
resp: { "conversation_id": "c_..." }
SSE 随后推出 run_started,带上这一轮的 run_id
2. GET /conversations/c_.../events
Accept: text/event-stream
(订阅这个会话的事件;连接掉了可以再挂,按 Last-Event-ID 续)
3. 用户看到方向不对
POST /conversations/c_.../commands
body: { "command_id": "cmd_...", "type": "redirect", "message": "先列风险" }
4. Agent 要删文件
SSE 推出:id: 1042
event: permission_required
data: {"run_id":"r_...","call_id":"call_12","tool":"rm"}
用户点允许
POST /conversations/c_.../approvals
body: { "command_id": "cmd_...", "call_id": "call_12", "allow": true }
第 3、第 4 步发生的时候,第 2 步那条 SSE 不必主动关掉。前端不用换会话 ID,也不用把已经渲染的半段字作废。后端只是往这个会话的事件日志里继续追加,再 fan-out 给所有订阅者。
画出来就是:
和 WebSocket 比,差别只剩一件事:回程消息不走这条流,走旁边一条短请求。事件仍然走原来那条长连接;连接自然断了,靠事件日志续,不靠「这条 TCP 还活着」。
这种「命令与事件分离、以应用层任务标识关联」的形状很常见。不过别把它等同于某个协议的永久固定设计。MCP 的历史有状态 Streamable HTTP 曾用 Mcp-Session-Id 串联 POST 和 GET/SSE;2026-07-28 修订↗已经拿掉协议层 session,改成每条请求自描述,跨调用状态靠工具参数里的显式 handle。OpenAI 的 Responses 支持 SSE 事件和 previous_response_id / Conversation,但它也不等于文中这套「自建事件日志 + 后台 run」。本文讨论的是应用层 Agent 的任务生命周期,不是在复述某一份协议。
先把三层 ID 拆开
一个 session_id 同时当对话、当一次 Loop、当 SSE 订阅范围,demo 里能跑,生产里会乱:用户在同一会话里发第二句,到底复用旧 run,还是新建 run?
拆成三层更干净:
| 标识 | 生命周期 | 作用 |
|---|---|---|
conversation_id | 一段用户任务 / 对话 | 上下文、成员、权限归属、SSE 订阅范围 |
run_id | 一次 Agent 执行尝试 | 取消、状态机、工具调用、审计 |
event_id | 单个事件 | 排序、去重、SSE 重放 |
done 不该是会话级终态。一次模型跑完,只是 run_completed / run_failed / run_cancelled,会话还可以收下一条命令。真正不再接受命令和订阅的,才发 conversation_closed。
对象再摊开一张表,后面的代码、竞态、幂等和多副本都对着它看:
| 对象 | 核心字段 | 典型状态 | 是否持久化 |
|---|---|---|---|
| Conversation | owner、messages | active / closed | 是 |
| Run | run_id、cancel token | queued / generating / waiting_approval / executing_tool / completed | 是 |
| Approval | call_id、参数摘要、审批人、过期时间 | pending / allowed / denied / expired | 是 |
| Event | event_id、run_id、type | token / tool / state | 至少短期持久化 |
Loop 还是那篇里的 Loop,只是多了两个口:
命令收件箱 ──→ 当前 run ──→ 事件日志 ──→ 每个 SSE 订阅者各自一份队列
↑
LLM / 工具
HTTP 负责往收件箱扔带 command_id 的命令,SSE 负责从事件日志拉给浏览器。连接可以断,run 可以继续;流可以重挂,事件可以从某个 event_id 接着播。生命周期绑在 conversation / run 上,不绑在某条 TCP 上。
两条路径具体怎么走
纠偏:打断当前 generation,但不主动拆流
用户插话时,模型多半正在吐 token。后端不能假装没看见,也不该把整个会话推倒重来。比较稳的做法是:
- HTTP 收到带
command_id的纠偏,按conversation_id + command_id去重后入队 - 当前 generation 必须是一个独立的 Task(或供应商提供的 cancellation handle)。Agent Loop 同时等两件事:模型流,和命令收件箱
- 收到
redirect:把当前run_id标成interrupting,取消模型请求,等这一轮清理完,再发generation_interrupted - 已经推出去的 token 不必撤回。前端打个「已打断」标记就行
- 在同一 conversation 里创建新的
run_id,把用户这句话追加进 messages,继续转 - 新的 token 仍然往原来那条事件日志上追加;还活着的 SSE 接着收,断了的按
Last-Event-ID续
前端看到的是连续的一条时间线:半段旧输出、一条系统提示、一段新输出。没有「旧连接死了,新连接生了」。
纠偏要落在run 边界上,不要幻想在某个 token 中间热切换模型内部状态。模型是无状态补全器,你能做的是:停掉这一轮请求,带着新的 messages 再问一次。
还有一句必须说谨慎:关掉客户端读流,不必然等于服务端停止推理或停止后台任务。 各家 SDK 的 close() / cancel() 能力不一样,有的只停本地迭代器,供应商那边还在烧 token。实际实现要以模型供应商是否支持服务端取消、你的 HTTP 客户端是否真的 abort 了请求为准。
取消模型,也不等于取消工具:
queued → generating → waiting_approval → executing_tool → completed
↓ ↓
rejected cancellation_requested
- 模型生成通常可中断
- 工具可能还没开始,取消很便宜
- 工具可能正在跑,但 API 不可中断
- 工具可能已经产生不可逆副作用(邮件发出去了、文件删了)
cancelled 不是「所有外部副作用都已撤销」。对不可撤销工具,应返回实际执行状态和补偿建议,而不是假装世界回到了点击之前。
授权:Loop 自己停住,SSE 继续心跳
授权跟纠偏相反:不是打断,是合法地卡住。
模型吐出 tool_calls
循环看见这是高风险工具
不执行,先往事件日志推:
id: 1042
event: permission_required
data: {"run_id":"r_...","call_id":"call_12","tool":"rm","args_digest":"..."}
然后 await pending_approvals[call_id] 这一把 Future
不要用「从一个共用 Queue 里取出、不是自己的就塞回去」来等审批。那会忙等、打乱顺序、让别的命令饥饿。正确的形状是:
pending_approvals: dict[call_id, Future]
发出 permission_required 时创建 Future;POST /approvals 只负责原子地完成它。同一 call_id 只允许批准或拒绝一次,重复提交返回已处理状态,过期的直接 expired。审批还要绑 actor_id + 参数摘要 + 过期时间,不能只认一个裸 ID。
SSE 此时还活着。与其推业务事件 ping,不如隔几秒推一行注释帧 : ping,前端不会当成业务事件误处理。用户点允许,POST 一次,对应 Future 被完成,Loop 被唤醒,执行工具,再把 tool_result 推出去。
整段时间里,事件通道没有因为授权被主动关掉。关的是 Loop 里的一把锁,不是那条流。高风险工具的参数也不要原样全量推给前端——令牌、隐私字段、内部路径留在服务端,卡片上只给用户决策需要的摘要。
最小概念骨架:命令、会话与订阅如何拆开
为突出职责边界,下面省略了完整鉴权、CSRF、超时回收、分布式恢复和供应商取消的细节;不能直接照搬到生产。它要钉死的是三件事:事件先落日志再广播、审批走 Future、Loop 同时听模型和命令。
事件日志,不是一只被抢着消费的 Queue
单个 asyncio.Queue 当「总线」会直接打脸后文的回放和多订阅:没人连时事件堆在里面还算运气好,两个标签页会把同一条事件抢走,断线后也无法按 ID 重放。
from dataclasses import dataclass, field
import asyncio, json, time, uuid
@dataclass
class Event:
event_id: int
name: str
data: dict
class Conversation:
def __init__(self, conversation_id: str, owner_id: str):
self.id = conversation_id
self.owner_id = owner_id
self.closed = False
self.event_log: list[Event] = []
self.subscribers: set[asyncio.Queue] = set()
self.inbox: asyncio.Queue = asyncio.Queue()
self.seen_commands: set[str] = set()
self.pending_approvals: dict[str, asyncio.Future] = {}
self.generation: asyncio.Task | None = None
self.current_run_id: str | None = None
def publish(self, name: str, data: dict) -> Event:
event = Event(len(self.event_log) + 1, name, data)
self.event_log.append(event)
for q in list(self.subscribers):
q.put_nowait(event)
return event
def subscribe(self, last_event_id: int) -> asyncio.Queue:
q: asyncio.Queue = asyncio.Queue()
for event in self.event_log:
if event.event_id > last_event_id:
q.put_nowait(event)
self.subscribers.add(q)
return q
publish() 先写日志,再 fan-out。每个 SSE 连接有自己的队列。新订阅按 Last-Event-ID 回放,再接实时事件。这才配叫事件通道。
生产里这份 event_log 要有上限或外置(Redis Streams、NATS JetStream、数据库事件表)。普通 NATS 核心、纯 Pub/Sub 没有回放,当不了这件事。
审批是一把按 call_id 索引的 Future
def request_approval(self, call_id: str, payload: dict) -> asyncio.Future:
loop = asyncio.get_running_loop()
fut = loop.create_future()
self.pending_approvals[call_id] = fut
self.publish("permission_required", payload)
return fut
def resolve_approval(self, call_id: str, allow: bool, actor_id: str) -> dict:
fut = self.pending_approvals.get(call_id)
if fut is None:
return {"status": "unknown_or_done"}
if fut.done():
return {"status": "already_resolved", "allow": fut.result()}
# 这里还应校验 actor_id、参数摘要、过期时间
fut.set_result(allow)
return {"status": "ok", "allow": allow}
POST /approvals 不再往 inbox 里塞一条「也许是审批」的消息,让 Loop 自己翻队列。谁发出 permission_required,谁 await 这把 Future。
Loop 必须同时等模型和收件箱
如果 run() 正在同步地跑完一轮 LLM 才回头看 inbox,纠偏 POST 入了队也只是「下一轮输入」,不是中断当前生成。
async def loop(self):
while not self.closed:
cmd = await self.inbox.get()
if cmd["type"] == "close":
break
await self.start_run(cmd)
while self.generation and not self.generation.done():
get_cmd = asyncio.create_task(self.inbox.get())
done, _ = await asyncio.wait(
{self.generation, get_cmd},
return_when=asyncio.FIRST_COMPLETED,
)
if self.generation in done:
get_cmd.cancel()
break
nxt = get_cmd.result()
if nxt["type"] in {"redirect", "close", "cancel"}:
self.generation.cancel()
try:
await self.generation
except asyncio.CancelledError:
pass
self.publish("generation_interrupted", {
"run_id": self.current_run_id,
})
if nxt["type"] == "redirect":
await self.start_run(nxt)
else:
self.closed = nxt["type"] == "close"
break
else:
await self.inbox.put(nxt)
async def start_run(self, cmd: dict):
self.current_run_id = str(uuid.uuid4())
self.publish("run_started", {"run_id": self.current_run_id})
self.generation = asyncio.create_task(self.generate(cmd))
async def generate(self, cmd: dict):
run_id = self.current_run_id
try:
# async for chunk in provider.stream(...):
# self.publish("token", {"run_id": run_id, "text": chunk})
# 需要授权时:
# allow = await self.request_approval(call_id, payload)
self.publish("run_completed", {"run_id": run_id})
except asyncio.CancelledError:
# 这里 close 的是你这边的 HTTP 读流。
# 供应商会不会停,取决于它支不支持服务端取消。
self.publish("run_cancelled", {"run_id": run_id})
raise
收到 redirect 的顺序是:标记中断 → 取消当前 Task → 等清理 → 发 generation_interrupted → 开新 run。asyncio.CancelledError 只保证你的协程停了,不保证对面模型或已经出手的工具也停了。
HTTP 接口:立刻还 ID,订阅只读事件
from fastapi import FastAPI, Header, HTTPException, Request
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
app = FastAPI()
conversations: dict[str, Conversation] = {}
def require_owner(conv: Conversation, request: Request) -> None:
user_id = getattr(request.state, "user_id", None)
if user_id is None or user_id != conv.owner_id:
raise HTTPException(403, "forbidden")
class CreateConversation(BaseModel):
command_id: str
message: str
class Command(BaseModel):
command_id: str
type: str = "redirect"
message: str
class Approval(BaseModel):
command_id: str
call_id: str
allow: bool
@app.post("/conversations")
async def create_conversation(body: CreateConversation, request: Request):
conv = Conversation(str(uuid.uuid4()), request.state.user_id)
conversations[conv.id] = conv
asyncio.create_task(conv.loop())
await conv.inbox.put({"type": "user", **body.model_dump()})
# run_id 由随后的 run_started 事件带出来,这里不抢一个可能还是 None 的值
return {"conversation_id": conv.id}
@app.post("/conversations/{conversation_id}/commands")
async def inject_command(conversation_id: str, body: Command, request: Request):
conv = conversations.get(conversation_id)
if not conv:
raise HTTPException(404)
require_owner(conv, request)
if body.command_id in conv.seen_commands:
return {"ok": True, "deduped": True}
conv.seen_commands.add(body.command_id)
await conv.inbox.put(body.model_dump())
return {"ok": True}
@app.post("/conversations/{conversation_id}/approvals")
async def approve(conversation_id: str, body: Approval, request: Request):
conv = conversations.get(conversation_id)
if not conv:
raise HTTPException(404)
require_owner(conv, request)
if body.command_id in conv.seen_commands:
return {"ok": True, "deduped": True}
conv.seen_commands.add(body.command_id)
return conv.resolve_approval(body.call_id, body.allow, request.state.user_id)
@app.get("/conversations/{conversation_id}/events")
async def events(
conversation_id: str,
request: Request,
last_event_id: str | None = Header(default=None, alias="Last-Event-ID"),
):
conv = conversations.get(conversation_id)
if not conv:
raise HTTPException(404)
require_owner(conv, request)
after = int(last_event_id) if last_event_id else 0
q = conv.subscribe(after)
async def gen():
try:
while not conv.closed:
try:
item = await asyncio.wait_for(q.get(), timeout=15)
except asyncio.TimeoutError:
yield ": ping\n\n"
continue
payload = json.dumps(item.data, ensure_ascii=False)
yield f"id: {item.event_id}\nevent: {item.name}\ndata: {payload}\n\n"
finally:
conv.subscribers.discard(q)
return StreamingResponse(
gen(),
media_type="text/event-stream; charset=utf-8",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
)
注意几个时间点和格式:
POST /conversations立刻返回 ID,不等模型说完。Agent 在后台自己转。GET /events不携带用户句子。它是订阅,不是提问。- 每条 SSE 都带
id:。浏览器原生EventSource重连时会自动带上Last-Event-ID,这才是第前面说的续播,不是口头承诺。 - 心跳用
: ping,不要发一个叫ping的业务事件。 - ID 用完整 UUID。截成 8 位十六进制只有 32 bit 随机空间,既容易撞,也不该当「谁知道这个 URL 谁就能审批」的凭据。随机 ID 也不是权限控制——每个命令和每次订阅都要校验当前用户是不是这个 conversation 的 owner。Cookie 鉴权还要补 CSRF。
前端仍然是两段互不打架的代码:
const { conversation_id } = await fetch("/conversations", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ command_id: crypto.randomUUID(), message }),
}).then((r) => r.json());
const es = new EventSource(`/conversations/${conversation_id}/events`);
es.addEventListener("token", (e) => appendText(JSON.parse(e.data).text));
es.addEventListener("permission_required", (e) => showApproval(JSON.parse(e.data)));
es.addEventListener("run_completed", () => markIdle());
es.addEventListener("generation_interrupted", () => markInterrupted());
await fetch(`/conversations/${conversation_id}/commands`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
command_id: crypto.randomUUID(),
type: "redirect",
message: "先列风险,先别改代码",
}),
});
EventSource 不会因为旁边多了几个 POST 就重连。这正是要的效果。run_completed 只表示这一轮闲下来,会话还可以继续接命令。
有一个竞态,必须先认
POST /conversations 一返回,Loop 可能已经开始吐 token;前端这时才去挂 SSE。中间那几十到几百毫秒的事件,如果只靠「当前这只 Queue」,新订阅者看不见。
这不是边角,是这套模型的第一道坑。三种收法,按常见程度排:
- 会话里留一段回放缓冲。事件带单调递增的
event_id,SSE 写出id:,客户端带Last-Event-ID。新挂上的从 0 或者上次看到的 id 重放。浏览器原生EventSource断线重连时会自动带上这个头,顺手就把「晚到订阅」和「中途掉线」一起解决了。 - 先建空会话,再挂流,再发第一句。三次请求,没有漏事件,但首字延迟多一拍。对「必须从第一个 token 就不断」的产品更合适。
- POST 的响应头里带
conversation_id,body 同时开始 SSE。这是老模型和会话模型的折中:第一句还是「请求即流」,后续插话改走独立 POST。能用,但第一条流的生命周期又和第一句绑在一起了,不如彻底拆干净。
生产里我更倾向 1。会话本来就要能重连,回放缓冲不是额外成本。上面那份骨架走的就是这条路。
几个不做就会在线上炸的点
取消当前生成要真的取消。 收件箱里有纠偏,只是「下一轮再用」不够。Loop 必须并发监听命令;取消时要 abort 你这边的请求,并核对供应商是否支持服务端停止。已经走到 executing_tool 的,按工具能不能停分别处理,不要对外宣称「已撤销」。
命令和审批都要幂等。 浏览器、SDK、网关都可能把同一个 POST 打两遍。command_id 去重纠偏和取消;call_id 保证同一工具调用只被批准或拒绝一次。第二次是 no-op,不要把工具执行两遍,也不要把同一句纠偏追加进 messages 两次。
粘性路由只能解决路由,不解决可靠性。 Loop、收件箱、事件日志默认都在进程内存里。换一个 Pod,ID 就成了 404。按 conversation_id 做粘性路由,能让请求回到同一个实例,却扛不住实例重启、发布迁移和进程崩溃。需要断线续传或可审计执行时,事件、审批状态和工具执行状态仍应外置并可恢复。能回放的是 Redis Streams、NATS JetStream、数据库事件表;普通消息总线不等于可重放事件日志。MCP 那篇里讲过有状态 Streamable HTTP 的粘性路由,问题同类,但别再把 MCP 当前协议层设计当成这套模型的依据。
SSE 要能闲着,也要承认它会断。 等授权可能是三十秒,也可能是三分钟。中间如果完全不推数据,网关、CDN、浏览器都可能默默拆连接。隔几秒推 : ping 比事后解释「卡片怎么没了」便宜。核心诉求不是「连接永远不掉」,而是:交互不主动拆流,掉了能按事件 ID 续。
会话要有死法。 用户关页、超时、主动 close,都得能把当前 run 停掉、把订阅队列摘掉、把未完成的 Future 取消掉。否则内存里全是半死不活的 asyncio.Task。run_completed 之后是留几分钟给前端拉回放,还是立刻 conversation_closed,产品自己定,但要显式定。
别用 GET 传第一句用户输入。 提问、纠偏、授权都是有副作用的命令,走 POST。SSE 的 GET 只读事件。缓存、重试、日志、权限校验都会简单很多。
那还要不要上 WebSocket
要,如果两头都在高频说话。
多人同时改一份 Agent 画布、协作光标、语音打断、客户端每 200ms 推一次本地状态,SSE 旁边再跟一串 POST 会显得傻,WSS 一条管子更干净。
但「偶尔纠偏 + 偶尔点授权」不是这种负载。它们是稀疏的、离散的、完全可以当成普通 API 的命令。SSE 的好处还在:那篇协议对比里写过——它就是 HTTP,浏览器帮你重连,网关当普通长请求处理,比维护一堆 WSS 心跳省心。重连能接上,前提是你真的写了 id: 和事件日志,不是只开着一条连接祈福。
所以选型可以收成一句:
事件为主、回程稀疏 → HTTP 命令 + 应用层会话 ID + SSE 订阅
两头持续互发 → WSS
问一次、不插话、不授权 → 请求体直接带流也行,别过度设计
第三条是早期 Chat 的形状,现在也没过时。过时的是:已经有了插话和授权,却还把流焊在第一句 POST 上。
小结
SSE 单向,并不妨碍 Agent 做人机协同。碍事的是把「用户这句话」当成「这条流的生命」。
换成应用层会话之后,关系变成:
HTTP → 往 conversation 里扔带 command_id 的命令
conversation → 养活上下文、当前 run、审批 Future、事件日志
SSE → 订阅事件日志;不因交互主动拆,断了按 event_id 续
用户改方向,是取消当前 run,在同一 conversation 里开一个新 run;Agent 要权限,是当前 run 停在工具门口,等一把按 call_id 索引的 Future。流不用为每一次人的介入陪葬,但也不要假装连接永不掉、取消等于世界回滚。
WebSocket 能做的,这套拆法在「稀疏回程」里同样能做,而且更贴现有的 HTTP 基础设施。等哪天回程也变成持续的、双向的、高频率的,再换 WSS 不迟。
