Skip to content

WorkspaceRealtimeSession 详解:单连接会话与 Registry / Broadcaster 的配合 ​

对象文件:backend/services/workspace_realtime/workspace_realtime_session.py。 解决两个问题:① 这个文件负责什么?② ConnectionRegistry 和 RealtimeBroadcaster 为什么有的地方一起调用、有的地方单独调用?

阅读建议:对着源码看,文中代码块标注了行号(如 :114 表示第 114 行)。


1. 这个文件负责什么(整体职责) ​

一句话:一个 WebSocket 连接 = 一个 WorkspaceRealtimeSession 实例。Session 绑定一个 user_id + 一个 workspace_id,负责这条连接从握手到断开的完整生命周期。

WS 端点(backend/api/realtime/workspace_realtime.py:54)非常薄,只做"Origin 校验 → 鉴权 → 构造 Session → await session.run()"。剩下所有事都在 Session 里。

Session 的 6 项职责 ​

职责对应方法干啥
握手鉴权authorize :102复用 JWT 解码 + require_view 成员校验,失败 close(1008)
连接注册run 开头 :114注册到本实例连接表(含连接数限制)
订阅广播run :123订阅该 workspace 的广播频道
在场态run :126-135拉在线快照单播 + 广播自己 join
上行事件路由_handle_upstream :257收到浏览器消息,按类型分发(presence/落库/透传)
可靠清理_teardown :344断开时退订、离开在场态、移出连接表

外加几个可靠性机制:_reauth_loop(周期重鉴权,被踢就断)、空闲超时 90s、单条消息容错。

Session 与周围组件的关系 ​

注意这张图的三条边:Session 用 Registry(只在连接生命周期)、用 Broadcaster(上行 publish + 订阅下行)、用 Presence(在场态)。Registry 和 Broadcaster 是 Session 依赖的两个独立工具,不是"绑定的套餐"——这是理解后续一切的关键。


2. 先建立全局认知:两条数据通路 ​

协同里所有实时数据流动,都可以归到两条通路:

  • 上行:Session → Broadcaster(只这俩,Registry 不参与)
  • 下行:Broadcaster → 各 Session(通过 _on_broadcast 回调,Registry 也不直接参与)

⚠️ 划重点:上下行都不直接遍历 ConnectionRegistry。这就是你看代码时觉得"有的地方没一起调用"的根源——它们本来就不该一起出现。


3. Registry 和 Broadcaster 到底各自管什么 ​

你的困惑本质上来自一个误解:以为它俩是"做同一件事的两个层次"。其实它们是两个正交的关注点:

ConnectionRegistryRealtimeBroadcaster
是什么本进程内存里的一张表跨进程的消息总线
数据结构dict[workspace_id → set[Session]]Redis Pub/Sub channel
作用范围单个实例(进程)整个集群(所有实例)
存的是Session 对象消息(handler 是订阅时注册的回调)
解决什么问题"我这个进程里,谁连着哪个 workspace?""一条消息怎么送到所有实例的所有 session?"
网络吗否,纯内存是,走 Redis

一个容易混淆的点:为什么"感觉有两份本实例连接数据"? ​

你看代码时可能会发现一个现象——本实例里连到某 workspace 的连接信息,似乎被记了两个地方:

  • ConnectionRegistry._connections[W] 装的是 Session 对象(用于:连接数限制、统计)
  • RedisBroadcaster._subscriptions[W] 装的是 handler 函数(用于:消息分发)

这不是冗余,是关注点分离:Registry 关心"连接对象"(限流时要数有多少个 Session),Broadcaster 关心"消息发到哪"(每个 Session 订阅时塞了一个回调进去)。理论上可以让 Broadcaster 直接持有 Session 并遍历调用 session.send(),但那样 Broadcaster 就要依赖 Session 类型,耦合更深;现在 Broadcaster 只认 handler 协议(任何 Callable[[str, dict], Any]),Session 只是个"能被回调"的订阅者,两边解耦。

重要发现:ConnectionRegistry 当前真实用途很窄 ​

我查了全仓调用点,ConnectionRegistry 在生产代码里只被调用了两处:

位置调用作用
workspace_realtime_session.py:114self._registry.add(self.workspace_id, self)建立连接时注册,顺带检查连接数上限
workspace_realtime_session.py:363self._registry.remove(self.workspace_id, self)断开时移出

而 ConnectionRegistry.online_user_ids() / sessions() 这两个方法目前没有任何生产代码调用(只在测试隔离时 _connections.clear())。

所以 ConnectionRegistry 当前真正干的事其实只有一件:连接数限制(add 时检查 MAX_CONNECTIONS_PER_WORKSPACE = 50,超限抛 ConnectionLimitError)。online_user_ids 是预留能力(未来做"本实例在线统计/排查"可用,现在还没接)。

这就解释了你的困惑:下行广播根本不遍历 Registry。那下行消息怎么发到本实例各连接的?答案在第 5 节。


4. 核心答疑:为什么它们不是 1:1 一起调用(场景对照表) ​

把 workspace_realtime_session.py 里所有涉及 Registry/Broadcaster 的调用点摊开,对照"谁被调、一起还是单独、为什么":

场景代码位置RegistryBroadcaster一起/单独为什么
连接建立run :114, :123add(限流)subscribe(_on_broadcast)一起既要点名(限流),也要开始收广播
连接断开_teardown :351, :363removeunsubscribe一起对称清理,两边都要注销
上行:用户操作要广播_persist_and_broadcast :331 / _handle_upstream :296❌ 不调publish单独(Broadcaster)发消息只要往 Redis 喊,不需要查本地连接表
下行:收到广播要发给浏览器_on_broadcast :242(被 Broadcaster 触发)❌ 不调(已是回调上下文)单独(Broadcaster)见第 5 节,分发靠订阅机制,不遍历 Registry
连接数超限run :114 抛异常add 抛 ConnectionLimitError❌ 不 subscribe单独(Registry)还没订阅就被拒了

一句话总结这个表:Registry 只在连接生命周期的边界(建立/断开)出现一次,负责"点名+限流";Broadcaster 贯穿所有消息流动(上行 publish、下行回调),负责"传消息"。它们只在生命周期边界处"碰巧一起出现",其余时候各管各的。


5. 核心困惑点:下行广播为什么不遍历 Registry?(本地扇出机制) ​

这是你最可能卡住的地方。先说你的直觉:

直觉里,下行广播应该是这样的—— Broadcaster 收到 Redis 消息 → 遍历 registry.sessions(workspace_id) → 对每个 Session send()。

但代码里根本找不到这行。因为实际机制不是这样。关键在 run() 里的这一句:

python
# backend/services/workspace_realtime/workspace_realtime_session.py:123
await self._broadcaster.subscribe(self.workspace_id, self._on_broadcast)
#                                              ^^^^^^^^^^^^^^^^^
#                                   每个 Session 把【自己的】_on_broadcast 注册进去

注意:订阅的粒度是"每个 Session 注册一个自己的 handler",不是"每个 workspace 注册一个总 handler"。所以 RedisBroadcaster 内部,_subscriptions[workspace_id] 这个 set 里,装的是本实例每个 Session 各自的 _on_broadcast。

当 Redis 一条消息到达,看 redis_broadcaster.py 的订阅循环:

python
# backend/services/workspace_realtime/redis_broadcaster.py:113(_pubsub_loop 内)
for handler in list(self._subscriptions.get(workspace_id, set())):
    result = handler(workspace_id, event)   # ← 遍历的是 Broadcaster 自己的 set,不是 Registry
    if asyncio.iscoroutine(result):
        await result

遍历的是 RedisBroadcaster._subscriptions 这个 set,不是 ConnectionRegistry._connections。每个 Session 的 _on_broadcast 各被触发一次,各自决定要不要 send 给自己的 ws:

python
# backend/services/workspace_realtime/workspace_realtime_session.py:242
async def _on_broadcast(self, workspace_id: str, event: dict[str, Any]) -> None:
    if workspace_id != self.workspace_id:     # 不是我关心的 workspace,跳过
        return
    sender = event.get("sender_session_id")
    if isinstance(sender, str) and sender == self._session_id:  # echo 排除(自己发的)
        return
    await self.send(event)                    # 发给【我自己绑的】那条 ws

完整机制图:

所以"遍历本实例所有连接"这件事,是被 Broadcaster 的订阅机制隐式完成的——因为每个 Session 建立时都自己 subscribe 了一个 handler,Broadcaster 自然就有了本实例所有连接的 handler 列表。ConnectionRegistry 在下行分发里完全不参与。

💡 这也是为什么你在 _on_broadcast、_persist_and_broadcast 里找不到 registry.sessions() 的调用——根本不需要。Registry 当前的活儿只有"连接数限制"那一处。


6. 连接生命周期时序图(标注每步谁被调) ​


7. 代码逐段对照(谁被调、为什么) ​

7.1 连接建立(run :109-184) ​

python
await self._websocket.accept()
try:
    self._registry.add(self.workspace_id, self)          # ← REGISTRY:点名 + 连接数限制
except Exception:                                         #    超限抛 ConnectionLimitError
    await self._websocket.close(code=1008); return
await self._broadcaster.subscribe(self.workspace_id, self._on_broadcast)  # ← BROADCASTER:开始收广播
# 注意:注册的是 self._on_broadcast(自己的回调),不是别的

snapshot = await self._presence.snapshot(workspace_id=self.workspace_id)  # ← PRESENCE:拉在线快照
if snapshot:
    await self.send({"type": EVENT_PRESENCE_STATE, "members": snapshot})  # 单播给新人
await self._presence.join(workspace_id=self.workspace_id, user_id=self.user_id)  # ← PRESENCE:广播 join

为什么一起调:这是连接生命周期边界——既要登记连接(限流),又要开启广播订阅。但注意它俩调用的目的不同,只是"恰好都在连接建立时"。

7.2 上行事件路由(_handle_upstream :257-306) ​

python
if event_type == EVENT_PRESENCE_CURSOR:
    await self._presence.update_cursor(...)      # ← PRESENCE(内部会 publish)
    return
if event_type in PERSISTED_EVENT_TYPES:
    await self._persist_and_broadcast(...)        # → 见 7.3
    return
if event_type in NOTIFY_EVENT_TYPES:
    await self._broadcaster.publish(self.workspace_id, {...})  # ← 只 BROADCASTER!不调 Registry
    return

为什么单独调 Broadcaster:用户操作要广播给别人,只要 publish(往 Redis 喊)。"发给谁"由 Redis 和各实例的订阅机制决定,发布方根本不需要知道本地有哪些连接,所以不碰 Registry。

7.3 落库 + 广播(_persist_and_broadcast :308-342) ​

python
async with AsyncSessionLocal() as db:
    handler = CanvasEventHandler(db, WorkspaceCanvasEventService(db))
    await handler.handle(..., auto_commit=False)   # 落库只 flush,不 commit
    await self._broadcaster.publish(self.workspace_id, {...})  # ← 只 BROADCASTER
    await db.commit()                              # 广播成功才 commit(失败回滚)

同样只调 Broadcaster,Registry 无关。

7.4 下行回调(_on_broadcast :242-255) ​

python
async def _on_broadcast(self, workspace_id, event):
    if workspace_id != self.workspace_id: return
    sender = event.get("sender_session_id")
    if isinstance(sender, str) and sender == self._session_id: return   # echo 排除
    await self.send(event)                          # 发给自己绑的 ws

这是 Broadcaster 回调 Session 的入口。不涉及 Registry——分发已经在 Broadcaster 内部遍历 _subscriptions 完成了。

7.5 清理(_teardown :344-368) ​

python
for step in (
    lambda: self._presence.leave(...),                              # PRESENCE
    lambda: self._broadcaster.unsubscribe(self.workspace_id, self._on_broadcast),  # BROADCASTER
):
    try: await step()
    except Exception: ...                          # 每步兜底,不阻断后续
self._registry.remove(self.workspace_id, self)     # ← REGISTRY:最后移出连接表

生命周期边界,对称清理:subscribe 对应 unsubscribe,add 对应 remove。这里两者都出现,因为是断开这个边界。注意三步各自 try 兜底,任一步(如 Redis 故障导致 unsubscribe 失败)不阻断后续——保证至少 Registry.remove(内存操作)一定执行,不留幽灵连接。


8. 一句话总结 + 常见误读纠正 ​

总结:ConnectionRegistry 是本进程的"连接点名册"(管限流,只在生命周期边界动),RealtimeBroadcaster 是跨集群的"消息总线"(管所有消息流动)。它们不是套餐,只在连接建立/断开时一起出现;消息上下行时只有 Broadcaster 在工作。

常见误读纠正:

误读正解
"Broadcaster 收到消息后会遍历 Registry 下发"不会。遍历的是 Broadcaster 自己的 _subscriptions(每个 Session 注册的 handler)
"Registry 和 Broadcaster._subscriptions 是冗余"不是。一个装 Session 对象(限流),一个装 handler(分发),关注点不同
"上行消息也要查 Registry 找到接收方"不要。publish 只管往 Redis 喊,接收方由订阅机制决定
"ConnectionRegistry 在协同里很核心"当前真实用途只有"连接数限制",online_user_ids/sessions 还没被业务调用(预留)
"subscribe 是订阅 workspace 级别的一个总入口"不是。是每个 Session 各 subscribe 一次自己的 _on_broadcast,所以 Broadcaster 才有本实例所有连接的 handler 列表

9. 相关文件 ​

文件职责
backend/services/workspace_realtime/workspace_realtime_session.py单连接会话(本文主角)
backend/services/workspace_realtime/connection_registry.py本实例连接表 + 连接数限制
backend/services/workspace_realtime/realtime_broadcaster.py广播协议 + 事件常量
backend/services/workspace_realtime/redis_broadcaster.pyRedis Pub/Sub 实现(_subscriptions 在这)
backend/api/realtime/workspace_realtime.pyWS 端点薄编排
docs/协作模式/2026-07-04-后端实现总结-agent产物同步与实时协同.md第 3.5 节讲跨实例广播分层,可与本文对照