Back to Journal
02 / Entry· 12 min read

别把 Agent 的 SSE 绑在一次提问上:用会话 ID 扛住中途插话和授权

Agent 不再是问一次答一次。用户要中途改方向,工具要等人授权。WebSocket 双向还好说,SSE 一旦跟某次提问焊死,流就得拆掉重开。把命令和订阅拆开:HTTP 拿应用层会话 ID,SSE 只负责挂在这个 ID 上。不因插话或授权主动拆流;连接自然断了,按事件 ID 续播。

🔊 系统朗读

上一篇把 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,交互模型非常干净:

text
用户提交一句话
POST /chat  (响应直接是 text/event-stream)
token 一段段推回来
模型说完,连接关掉

这跟普通补全 API 几乎是同一个形状,只是把整段 JSON 换成了 SSE。前端用 fetch 读流,或者用 EventSource 挂上去,把 data: 拼到对话框里。单向,刚好。

那时候 Agent Loop 也短。Agent Loop 那篇里写过:模型点名,循环执行,再问,停机。用户通常只在开头开口,中间最多看一眼工具调用卡片,并不参与。

所以「一次提问 = 一条 SSE」成了很多项目的默认实现。它简单,本地 demo 也好看。

然后产品开始要两件事。

中途插话和授权,把单向流逼到墙角

第一件是纠偏。模型已经开始写方案了,用户一眼看出方向不对:「别写 React,用 Vue。」「先别改代码,先列风险。」如果等它说完再发下一条,token 白烧,体验也像在跟一个听不见人的人说话。

第二件是授权。Agent 要执行有副作用的工具:覆盖文件、发邮件、调生产接口、读隐私数据。这时候循环必须停住,把「我准备干什么」推给前端,等人点允许或拒绝,再继续。这就是人机协同(HITL)最常见的一种。

WebSocket 处理这两件事更直接:

text
同一条 WSS
  ├─ 服务端:token / 工具状态 / 等授权
  └─ 客户端:纠偏 / 同意 / 拒绝

双向消息走同一条连接,少一次「再开一条流」的心智负担。但会话是否持续、断线后如何恢复,仍然是应用层状态管理问题。WebSocket 断了,如果 run、审批、事件还只活在内存里,照样要自己恢复——它并不比 SSE 更「自带状态」。

SSE 没有这条回程。规范就是服务端往下推,浏览器不能在同一条 EventSource 上再塞一个 body 上去。于是很多实现会做成这样:

text
用户点「允许删除」
前端 abort 当前 EventSource
再 POST 一次,带上 approval
后端新开一条 SSE
前端重新订阅、重新拼状态

能跑。但有三处硌手:

  1. 观感断了。流式输出最怕的就是中间闪一下。重连、重挂、重新对齐光标,用户会觉得系统卡了或者重来了。
  2. run 的生命周期被绑死在连接上。连接一断,后端如果把这次 Agent 循环一起拆掉,授权到来时那个「等在工具门口」的协程已经没了;如果不断,你又得自己发明一套「断流但 run 还活着」的状态机——而这套状态机,本来就不该依赖连接。
  3. 前端状态难交接。旧流里已经推了半段文本、两次工具调用、一张权限卡片。新流从哪接着播?重放全部?只播增量?漏一条事件,UI 就会和后端各活各的。

所以别被「SSE 不能回传」吓到。浏览器本来就能同时开很多 HTTP 请求。SSE 只是不能在这一条连接上反向说话,并不禁止你再发一个 POST。别扭的是那个默认假设:这条流是为「这一句用户输入」而生的,输入没了,流就该死。

该拆开的是命令和订阅,不是换协议

把一次 Agent 任务看成两个通道:

通道方向干什么用什么
命令通道客户端 → 服务端提问、纠偏、授权、取消普通 HTTP POST
事件通道服务端 → 客户端token、工具开始/结束、等授权、run 结束SSE

它们共享的不是某次 HTTP 请求,而是一个应用层交互会话。这里的 conversation_id 不是登录 Cookie 里的 Session,也不是一次模型请求。它只是把「这一次用户任务」的命令和事件对上号。

交互变成:

text
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 给所有订阅者。

画出来就是:

事件日志Agent RunHTTP 命令前端事件日志Agent RunHTTP 命令前端POST /conversations 第一句conversation_id + run_idGET /conversations/id/eventstoken / tool_call往下推POST /commands redirect取消当前 generationgeneration_interrupted开启新 run继续推事件permission_requiredPOST /approvals完成对应 Futuretool_result / 后续 tokenrun_completed

和 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

对象再摊开一张表,后面的代码、竞态、幂等和多副本都对着它看:

对象核心字段典型状态是否持久化
Conversationowner、messagesactive / closed
Runrun_id、cancel tokenqueued / generating / waiting_approval / executing_tool / completed
Approvalcall_id、参数摘要、审批人、过期时间pending / allowed / denied / expired
Eventevent_id、run_id、typetoken / tool / state至少短期持久化

Loop 还是那篇里的 Loop,只是多了两个口:

text
命令收件箱 ──→  当前 run ──→ 事件日志 ──→ 每个 SSE 订阅者各自一份队列
                 ↑
            LLM / 工具

HTTP 负责往收件箱扔带 command_id 的命令,SSE 负责从事件日志拉给浏览器。连接可以断,run 可以继续;流可以重挂,事件可以从某个 event_id 接着播。生命周期绑在 conversation / run 上,不绑在某条 TCP 上。

两条路径具体怎么走

纠偏:打断当前 generation,但不主动拆流

用户插话时,模型多半正在吐 token。后端不能假装没看见,也不该把整个会话推倒重来。比较稳的做法是:

  1. HTTP 收到带 command_id 的纠偏,按 conversation_id + command_id 去重后入队
  2. 当前 generation 必须是一个独立的 Task(或供应商提供的 cancellation handle)。Agent Loop 同时等两件事:模型流,和命令收件箱
  3. 收到 redirect:把当前 run_id 标成 interrupting,取消模型请求,等这一轮清理完,再发 generation_interrupted
  4. 已经推出去的 token 不必撤回。前端打个「已打断」标记就行
  5. 在同一 conversation 里创建新的 run_id,把用户这句话追加进 messages,继续转
  6. 新的 token 仍然往原来那条事件日志上追加;还活着的 SSE 接着收,断了的按 Last-Event-ID

前端看到的是连续的一条时间线:半段旧输出、一条系统提示、一段新输出。没有「旧连接死了,新连接生了」。

纠偏要落在run 边界上,不要幻想在某个 token 中间热切换模型内部状态。模型是无状态补全器,你能做的是:停掉这一轮请求,带着新的 messages 再问一次。

还有一句必须说谨慎:关掉客户端读流,不必然等于服务端停止推理或停止后台任务。 各家 SDK 的 close() / cancel() 能力不一样,有的只停本地迭代器,供应商那边还在烧 token。实际实现要以模型供应商是否支持服务端取消、你的 HTTP 客户端是否真的 abort 了请求为准。

取消模型,也不等于取消工具:

text
queued → generating → waiting_approval → executing_tool → completed
                         ↓                    ↓
                      rejected         cancellation_requested
  • 模型生成通常可中断
  • 工具可能还没开始,取消很便宜
  • 工具可能正在跑,但 API 不可中断
  • 工具可能已经产生不可逆副作用(邮件发出去了、文件删了)

cancelled 不是「所有外部副作用都已撤销」。对不可撤销工具,应返回实际执行状态和补偿建议,而不是假装世界回到了点击之前。

授权:Loop 自己停住,SSE 继续心跳

授权跟纠偏相反:不是打断,是合法地卡住

text
模型吐出 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 里取出、不是自己的就塞回去」来等审批。那会忙等、打乱顺序、让别的命令饥饿。正确的形状是:

text
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 重放。

python
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

python
    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 入了队也只是「下一轮输入」,不是中断当前生成。

python
    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,订阅只读事件

python
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。

前端仍然是两段互不打架的代码:

javascript
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」,新订阅者看不见。

这不是边角,是这套模型的第一道坑。三种收法,按常见程度排:

  1. 会话里留一段回放缓冲。事件带单调递增的 event_id,SSE 写出 id:,客户端带 Last-Event-ID。新挂上的从 0 或者上次看到的 id 重放。浏览器原生 EventSource 断线重连时会自动带上这个头,顺手就把「晚到订阅」和「中途掉线」一起解决了。
  2. 先建空会话,再挂流,再发第一句。三次请求,没有漏事件,但首字延迟多一拍。对「必须从第一个 token 就不断」的产品更合适。
  3. 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.Taskrun_completed 之后是留几分钟给前端拉回放,还是立刻 conversation_closed,产品自己定,但要显式定。

别用 GET 传第一句用户输入。 提问、纠偏、授权都是有副作用的命令,走 POST。SSE 的 GET 只读事件。缓存、重试、日志、权限校验都会简单很多。

那还要不要上 WebSocket

要,如果两头都在高频说话。

多人同时改一份 Agent 画布、协作光标、语音打断、客户端每 200ms 推一次本地状态,SSE 旁边再跟一串 POST 会显得傻,WSS 一条管子更干净。

但「偶尔纠偏 + 偶尔点授权」不是这种负载。它们是稀疏的、离散的、完全可以当成普通 API 的命令。SSE 的好处还在:那篇协议对比里写过——它就是 HTTP,浏览器帮你重连,网关当普通长请求处理,比维护一堆 WSS 心跳省心。重连能接上,前提是你真的写了 id: 和事件日志,不是只开着一条连接祈福。

所以选型可以收成一句:

text
事件为主、回程稀疏     →  HTTP 命令 + 应用层会话 ID + SSE 订阅
两头持续互发           →  WSS
问一次、不插话、不授权 →  请求体直接带流也行,别过度设计

第三条是早期 Chat 的形状,现在也没过时。过时的是:已经有了插话和授权,却还把流焊在第一句 POST 上。

小结

SSE 单向,并不妨碍 Agent 做人机协同。碍事的是把「用户这句话」当成「这条流的生命」。

换成应用层会话之后,关系变成:

text
HTTP           →  往 conversation 里扔带 command_id 的命令
conversation   →  养活上下文、当前 run、审批 Future、事件日志
SSE            →  订阅事件日志;不因交互主动拆,断了按 event_id 续

用户改方向,是取消当前 run,在同一 conversation 里开一个新 run;Agent 要权限,是当前 run 停在工具门口,等一把按 call_id 索引的 Future。流不用为每一次人的介入陪葬,但也不要假装连接永不掉、取消等于世界回滚。

WebSocket 能做的,这套拆法在「稀疏回程」里同样能做,而且更贴现有的 HTTP 基础设施。等哪天回程也变成持续的、双向的、高频率的,再换 WSS 不迟。

分享
← 返回博客列表
🎁 有邀请福利哦,点击查看
🎁