Appearance
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 到底各自管什么
你的困惑本质上来自一个误解:以为它俩是"做同一件事的两个层次"。其实它们是两个正交的关注点:
| ConnectionRegistry | RealtimeBroadcaster | |
|---|---|---|
| 是什么 | 本进程内存里的一张表 | 跨进程的消息总线 |
| 数据结构 | 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:114 | self._registry.add(self.workspace_id, self) | 建立连接时注册,顺带检查连接数上限 |
workspace_realtime_session.py:363 | self._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 的调用点摊开,对照"谁被调、一起还是单独、为什么":
| 场景 | 代码位置 | Registry | Broadcaster | 一起/单独 | 为什么 |
|---|---|---|---|---|---|
| 连接建立 | run :114, :123 | add(限流) | subscribe(_on_broadcast) | 一起 | 既要点名(限流),也要开始收广播 |
| 连接断开 | _teardown :351, :363 | remove | unsubscribe | 一起 | 对称清理,两边都要注销 |
| 上行:用户操作要广播 | _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)→ 对每个 Sessionsend()。
但代码里根本找不到这行。因为实际机制不是这样。关键在 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.py | Redis Pub/Sub 实现(_subscriptions 在这) |
backend/api/realtime/workspace_realtime.py | WS 端点薄编排 |
docs/协作模式/2026-07-04-后端实现总结-agent产物同步与实时协同.md | 第 3.5 节讲跨实例广播分层,可与本文对照 |