Appearance
协作模式后端实现总结:Agent 产物同步与实时协同
本文是对协作模式后端已落地代码的总结解读,聚焦两件事:
- Agent 产物如何同步为画布产物(
WorkspaceOutput)。- WebSocket + Redis Pub/Sub 如何实现实时协作。
上游设计文档:
docs/协作模式/2026-06-30-协作模式后端重构与实现设计.md。 本文不重复设计取舍,只对照实现讲清楚核心逻辑流程 + 核心代码片段 + 架构图。文件路径均为项目根相对路径(如
backend/services/...)。
1. 总体架构
协作模式后端由两个正交子系统组成:
- 产物同步子系统(
backend/services/workspace/):把 Agent 在会话里生成的媒体结果落成画布可定位的WorkspaceOutput。 - 实时协同子系统(
backend/services/workspace_realtime/):在多个协作者之间同步在场态、画布操作与产物变更。
二者通过一个共享抽象 RealtimeBroadcaster 解耦:产物同步完成后广播事件,实时协同子系统负责分发到每个在线连接。
模块结构图
目录职责约定(取自 docs/协作模式/2026-06-30-协作模式后端重构与实现设计.md):
| 目录 | 承载 |
|---|---|
services/workspace/ | 工作空间领域主流程与编排(业务真相) |
services/workspace/output_sync/ | 产物同步来源的可替换策略族(变化点 1) |
services/workspace_realtime/ | 实时协同基础能力,正交于领域主流程(变化点 2:传输实现) |
2. Agent 产物如何同步为画布产物
2.1 设计思路:依赖倒置 + 策略族
核心问题是:Agent 在会话里跑出图片/视频/音频,如何让它们自动出现在协作画布上,且不把 stream 服务与 workspace 业务耦合。
采用两层抽象:
- 依赖倒置(Observer 协议):
image/video/omni/chat四个 stream 服务(加上voice_service、action_mimic_service)在构造时注入一个WorkspaceOutputObserver,默认是NoopWorkspaceOutputObserver。stream 服务只认协议,不知道 workspace 存在与否。 - 策略族(产物同步来源):真正落库的逻辑封装成
OutputSyncStrategy,目前 AI 来源由ArtifactLinkSyncStrategy实现,上传/手动文本为 no-op(走 REST)。Observer 只持有一个Resolver,不感知具体来源。
workspace_id 的解析放在 Observer 内部(resolve_workspace_id),stream 服务只需传 conv_id,非 workspace 会话直接短路返回。
2.2 类交互图
2.3 核心组件详解
Observer 协议与 Default 实现
backend/services/workspace/workspace_output_observer.py 定义四个钩子,覆盖三种生成节奏:
| 钩子 | 触发时机 | 作用 |
|---|---|---|
on_generation_started | stream 开始生成、已拿到 assistant_message_id | 建 pending 占位产物,广播 output:add |
on_generation_completed | stream 生成结束、final_result/artifact_links 已定 | 更新 payload/status,广播 output:update |
on_video_task_submitted | 视频异步任务提交成功 | 为每个 video_task 建独立 pending 产物(带 task_id) |
on_audio_task_submitted | 长文本配音异步提交成功 | 建 pending audio 产物(带 task_id) |
视频/长文本配音是异步任务,stream 结束时 video_url 仍为空,所以不能走 on_generation_completed(会被误标 completed),而是在任务提交阶段单独建 pending 占位,前端轮询到终态后再回写 url(见 2.7)。
DefaultWorkspaceOutputObserver 每个钩子的骨架一致:开新 AsyncSessionLocal → 解析 workspace_id(非 workspace 会话 return)→ 走 ArtifactLinkSyncStrategy 落库 → 通过 _publish 广播。落库异常只 WARNING 不冒泡,保证 stream 主流程不受拖累。
产物同步策略族
backend/services/workspace/output_sync/output_sync_strategy_resolver.py 持有三套策略,未知来源回退到 AI 策略。Observer 直接通过 resolver.artifact_link_strategy 拿到 AI 策略使用。
ArtifactLinkSyncStrategy(backend/services/workspace/output_sync/artifact_link_sync_strategy.py)的核心职责:
sync_started:按(workspace_id, message_id)幂等——若已有同 message_id 产物则直接返回,否则建 pending 占位。sync_completed:按产物类型从final_result提取媒体 payload(图片走image_urls、视频走video_tasks[0].video_url、音频走audio_url)。多图场景首张写入占位、其余追加为独立产物;无媒体内容(纯文本/代码沙箱artifact_links)时软删占位,避免画布残留空卡片。- 画布布局:按已存在产物数量错峰摆放(百分比坐标),并根据图片宽高比修正画布框比例(
_completed_canvas_config)。
组合根
backend/services/workspace/observer_factory.py 提供 get_workspace_output_observer() 生产单例:构造 DefaultWorkspaceOutputObserver(broadcaster=get_realtime_broadcaster());若 broadcaster 解析失败(测试环境未初始化 Redis)则回退 Noop。各 stream 服务的 Cls.get_xxx_stream_service() 单例工厂在构造时注入该 observer。
2.4 完整时序:同步生成(以图片为例)
注:图片 stream 的开始/完成挂接点分别在
backend/services/core/image_stream_service.py:182与:411;四个 stream 服务挂接点一致(见 2.6 表格)。
2.5 异步生成时序(视频/长文本配音)
视频与配音提交即返回 task_id,真正的 url 在前端轮询到 done 后才有。因此提交阶段建 pending 占位、完成阶段不落库,回写交给 REST。
2.6 产物状态机
软删使用
deleted_at,查询一律deleted_at IS NULL过滤。
2.7 REST 辅助产物接口
AI 自动同步是主路径,REST 产物接口承担"手动落产物 + 异步任务回写 + 布局/内容编辑 + 软删"。集中在 backend/api/workspace.py 与 backend/services/workspace/workspace_output_service.py:
| 方法 | 路由 | 说明 |
|---|---|---|
| GET | /api/workspaces/{id}/outputs | 拉取画布全部未软删产物(进入画布首屏) |
| POST | /api/workspaces/{id}/outputs | 手动落产物(上传/文本),AI 产物不走这里 |
| PUT | /api/workspaces/{id}/outputs/{output_id} | 更新布局/内容/状态;视频/配音轮询回写 url 走这里 |
| DELETE | /api/workspaces/{id}/outputs/{output_id} | 软删产物(限作者或 owner) |
权限模型(workspace_output_service.py):布局 canvas_config 人人可拖(协作刚需);但修改内容 payload/status 限产物作者或项目 owner,防止成员篡改他人产物。_require_owner_or_author 集中校验。
说明:
ArtifactLinkSyncStrategy.update_video_output/update_audio_output两个方法当前在策略类中定义但未见调用点,异步任务回写实际由前端轮询到终态后调用通用 PUT 产物接口完成。文档如实标注,避免误解。
2.8 生成侧挂接点一览
stream 服务对 workspace 零反向耦合,统一通过构造参数注入 observer(默认 Noop):
| 文件 | 挂接点行号 | 钩子 |
|---|---|---|
backend/services/core/image_stream_service.py | 182 / 411 | started / completed |
backend/services/core/video_stream_service.py | 643 / 660 | completed / video_task_submitted |
backend/services/core/omni_stream_service.py | 205 / 393 / 410 | started / completed / video_task_submitted |
backend/services/core/chat_stream_service.py | 406 | completed |
backend/services/voice_service.py | 349 / 503 | completed / audio_task_submitted |
backend/services/action_mimic_service.py | 215 / 351 | video_task_submitted |
3. WebSocket + Redis Pub/Sub 实时协同
3.1 通信模型与设计目标
- 通信模型:WebSocket 长连接 + Redis Pub/Sub 跨实例广播;画布操作"服务端权威落库 + 广播确认",无字符级合并(不是 CRDT)。
- 设计目标:
- 单实例能跑,多实例也能扇出(broadcaster 可替换)。
- 业务代码只依赖
RealtimeBroadcaster协议,传输替换无感。 - 测试不依赖 Redis(注入
FakeBroadcaster)。
3.2 类交互图
3.3 连接生命周期
WS 端点 backend/api/realtime/workspace_realtime.py 是薄编排层,只做四件事:Origin 校验 → 鉴权 → 构造 Session → 运行。
握手鉴权复用 backend/core/jwt.py:decode_access_token,经 WorkspaceAccessControl.require_view 校验成员级可见性,失败 close(1008)。
3.4 事件类型与路由分发
事件常量集中在 backend/services/workspace_realtime/realtime_broadcaster.py,按"服务端要不要落库"分三类。WorkspaceRealtimeSession._handle_upstream 是唯一的上行事件路由:
| 分类 | 事件类型 | 处理方式 |
|---|---|---|
| 静默 | presence:ping | 客户端 25s 心跳保活,服务端忽略不广播 |
| 在场态 | presence:cursor / presence:select | 写 Redis presence hash + 广播增量,不入库 |
| 权威落库 | canvas:event / output:update | 服务端校验+落库(含产物布局更新),广播确认 |
| 通知透传 | output:add/output:delete/output:content/annotation:* | 前端 REST 落库后发起,服务端不落库,原样转发(带 sender_session_id) |
3.5 广播分层:单实例内存 + 跨实例 Redis
这是实时协同的核心机制。广播能力分两层:
ConnectionRegistry(单实例内):维护workspace_id → set[Session]的内存字典,记录"本进程内谁连着哪个 workspace"。RedisBroadcaster(跨实例):每个 workspace 一个订阅协程,订阅 Redis channelworkspace:realtime:{workspace_id},收到消息后分发给本实例注册的本地 handler。
关键工程细节(backend/services/workspace_realtime/redis_broadcaster.py):
- publish 与 pubsub 用不同连接池。
pubsub()会独占一条连接直到关闭,若与 publish 共享max_connections=50的池,活跃 workspace 数达上限后会耗尽池导致 publish 失败。所以 pubsub 用独立池(_make_pubsub_redis,max_connections=64)。 - 每个 workspace 一个订阅协程,首条订阅时创建,最后一个 handler 退订时取消。
- 重连退避:
_pubsub_loop外层while True,Redis 抖动/连接重置后sleep(_RECONNECT_BACKOFF_SECONDS)重订阅,保证跨实例广播在该实例持续可用(不"失明");仅unsubscribe触发的CancelledError才终止。 - 严格模式不降级:生产组合根
get_realtime_broadcaster()直接构造RedisBroadcaster,Redis 未初始化即抛错。
3.6 在场态(Presence)
backend/services/workspace_realtime/workspace_presence_service.py 把"谁在线、光标在哪、选了哪个产物"存到 Redis hash:
- key:
workspace:presence:{workspace_id},field=user_id,value=JSON{user_id, cursor, selection}。 - TTL 120s,每次写入(cursor/select/join)都续期,等价于心跳续期;超过 TTL 无活动即视为离线,被
snapshot排除。 snapshot读取全量成员供新连接初始同步:新人一进画布就能看到"已在画布且静止"的协作者光标/选区,无需等他们再动。- join/leave/cursor/select 同时经 broadcaster 广播增量事件。
- 跨实例共享:多实例读同一个 Redis hash,在线成员全局一致。
3.7 关键时序
3.7.1 画布操作:权威落库 + 广播确认(落库广播原子化)
这是协作画布防漂移的关键设计(M7):落库只 flush 不 commit,广播成功后再 commit,广播失败则整体回滚。
并发拖拽同一产物时,workspace_canvas_event_service.py 用 SELECT ... FOR UPDATE 加行锁,防止 last-write-wins 丢更新。
3.7.2 产物通知:透传 + echo 排除
新增/删除产物、批注变更由前端先 REST 落库,再经 WS 广播给其他成员。服务端不落库只透传,并附 sender_session_id 让发起方不收到自己刚发的事件(画布操作是乐观更新,收不到自己的确认无影响)。
3.7.3 AI 产物同步 → 实时广播(子系统衔接)
DefaultWorkspaceOutputObserver 落库后调用 broadcaster.publish,把 output:add/output:update 事件扇出到所有在线成员(包括异实例),这就是 2.4 时序里"其他成员收到更新"的实现路径。
3.8 核心代码片段
Redis Pub/Sub 订阅循环(redis_broadcaster.py)
python
# backend/services/workspace_realtime/redis_broadcaster.py
async def _pubsub_loop(self, workspace_id: str) -> None:
"""订阅 Redis channel,把消息分发给本地 handler。外层重连循环。"""
while True:
try:
async with self._pubsub_redis.pubsub(ignore_subscribe_messages=True) as pubsub:
await pubsub.subscribe(_channel(workspace_id))
while True:
message = await pubsub.get_message(timeout=1.0)
if message is None or message.get("type") != "message":
continue
event = self._decode(message.get("data"))
if event is None:
continue
for handler in list(self._subscriptions.get(workspace_id, set())):
result = handler(workspace_id, event)
if asyncio.iscoroutine(result):
await result
except asyncio.CancelledError:
raise
except Exception:
logger.error("workspace.realtime.redis pubsub error, will retry workspace_id=%s",
workspace_id, exc_info=True)
await asyncio.sleep(_RECONNECT_BACKOFF_SECONDS) # 退避重连上行事件路由(workspace_realtime_session.py)
python
# backend/services/workspace_realtime/workspace_realtime_session.py
async def _handle_upstream(self, event: dict[str, Any]) -> None:
event_type = str(event.get("type") or "")
payload = event.get("payload") if isinstance(event.get("payload"), dict) else {}
output_id = event.get("output_id") if isinstance(event.get("output_id"), int) else payload.get("output_id")
if event_type in SILENT_EVENT_TYPES: # presence:ping 静默
return
if event_type == EVENT_PRESENCE_CURSOR: # 光标聚合,不入库
await self._presence.update_cursor(workspace_id=self.workspace_id,
user_id=self.user_id,
x=_safe_float(payload.get("x")),
y=_safe_float(payload.get("y")))
return
if event_type == EVENT_PRESENCE_SELECT: # 选区聚合,不入库
await self._presence.update_selection(...)
return
if event_type in PERSISTED_EVENT_TYPES: # 画布操作:权威落库 + 广播确认
await self._persist_and_broadcast(event_type=event_type, payload=payload, output_id=...)
return
if event_type in NOTIFY_EVENT_TYPES: # 通知类:原样透传 + echo 排除
await self._broadcaster.publish(self.workspace_id, {
**event, "workspace_id": self.workspace_id,
"user_id": self.user_id, "sender_session_id": self._session_id,
})
return
# 其余未知事件一律丢弃落库与广播原子化(workspace_realtime_session.py)
python
# backend/services/workspace_realtime/workspace_realtime_session.py
async def _persist_and_broadcast(self, *, event_type, payload, output_id) -> None:
"""落库只 flush 不 commit,广播成功后再 commit;广播失败 async with 退出自动回滚。"""
async with AsyncSessionLocal() as db:
handler = CanvasEventHandler(db, WorkspaceCanvasEventService(db))
await handler.handle(user_id=self.user_id, workspace_id=self.workspace_id,
event_type=event_type, payload=payload, output_id=output_id,
auto_commit=False) # 延迟提交
await self._broadcaster.publish(self.workspace_id, {
"sender_session_id": self._session_id, "type": event_type,
"workspace_id": self.workspace_id, "user_id": self.user_id,
**({"output_id": output_id} if output_id is not None else {}),
**({"payload": payload} if payload else {}),
})
await db.commit() # 广播成功后才提交Observer 落库 + 广播(workspace_output_observer.py)
python
# backend/services/workspace/workspace_output_observer.py
async def on_generation_started(self, *, conv_id, message_id, output_type, hint=None) -> None:
async with AsyncSessionLocal() as db:
workspace_id = await resolve_workspace_id(db, conv_id) # 非 workspace 会话 return
if workspace_id is None:
return
try:
resolver = OutputSyncStrategyResolver(db)
output = await resolver.artifact_link_strategy.sync_started(
workspace_id=workspace_id, conv_id=conv_id, message_id=message_id,
output_type=output_type, hint=hint,
)
await db.commit()
except Exception:
logger.warning("workspace.output.sync_started failed conv_id=%s message_id=%s",
conv_id, message_id, exc_info=True)
return
if output is not None:
await self._publish(workspace_id, {"type": "output:add", "workspace_id": workspace_id, "output": output})3.9 可靠性设计一览
实时协同对连接泄漏、状态漂移敏感,代码里有一组防御性设计:
| 机制 | 位置 | 作用 |
|---|---|---|
| 空闲超时 90s | _IDLE_TIMEOUT_SECONDS + wait_for(receive_json) | 客户端半开(NAT 黑洞)时不致连接与订阅永久泄漏 |
| 周期重鉴权 60s | _reauth_loop | 握手后被踢出 workspace 仍能及时切断连接 |
| 单条消息容错 | _handle_upstream try/except | 业务异常只跳过本条,不冒泡关闭整条连接 |
| 落库广播原子化 | _persist_and_broadcast flush→publish→commit | 广播失败回滚,避免 DB 已变但他人收不到导致漂移 |
| 连接数上限 50/实例/workspace | ConnectionRegistry.add | 防刷连接耗尽内存与 Redis 订阅 |
| echo 排除 | sender_session_id | 发起方不收到自己乐观更新的事件回环 |
| pubsub 重连退避 | _pubsub_loop | Redis 抖动后自动重订阅,不"失明" |
| teardown 多步兜底 | _teardown 逐 step try/except | 任一步异常不阻断后续清理,避免幽灵 session |
| Origin 校验 | _is_origin_allowed | 防跨站 WebSocket 劫持(CSWSH) |
4. 关键设计点速查
| 关注点 | 选型 | 理由 |
|---|---|---|
| 产物同步挂接 | Observer 协议 + 依赖注入 | stream 服务与 workspace 业务零反向耦合 |
| 产物同步来源 | 策略族 + Resolver | 开闭原则,未来批注转产物/模板插入可加策略 |
| workspace 归属判定 | Observer 内部 resolve_workspace_id | 归属判定不散落到各 stream 服务 |
| 产物幂等 | (workspace_id, message_id) 找最早占位 | 一轮多图不重复建卡 |
| 视频异步 | 提交建 pending + REST 回写 | stream 结束时 url 为空,不能走 completed |
| 广播传输 | 协议 + 内存/Redis 双实现 | 单实例可跑,多实例可扇出,测试可注入 Fake |
| 画布一致性 | 服务端权威落库 + 广播确认 | 避免 CRDT 复杂度,以 DB 为真相源 |
| 跨实例在场态 | Redis hash + TTL | 多实例读同一 hash,在线状态全局一致 |
5. 文件索引
产物同步子系统
| 文件 | 职责 |
|---|---|
backend/services/workspace/workspace_output_observer.py | Observer 协议 + Noop/Default + resolve_workspace_id |
backend/services/workspace/observer_factory.py | 生产 observer 单例(依赖 broadcaster) |
backend/services/workspace/output_sync/output_sync_strategy.py | 策略协议 |
backend/services/workspace/output_sync/artifact_link_sync_strategy.py | AI 产物落库核心(pending/多图/视频/配音/画布布局) |
backend/services/workspace/output_sync/upload_sync_strategy.py | 上传策略(no-op,走 REST) |
backend/services/workspace/output_sync/manual_text_sync_strategy.py | 手动文本策略(no-op) |
backend/services/workspace/output_sync/output_sync_strategy_resolver.py | 策略解析器 |
backend/services/workspace/workspace_output_service.py | REST 手动产物(查询/增改/软删/权限) |
backend/models/workspace_output.py | 画布产物模型 |
实时协同子系统
| 文件 | 职责 |
|---|---|
backend/api/realtime/workspace_realtime.py | WS 端点薄编排(Origin/鉴权/构造 Session) |
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 实现(独立池/重连) |
backend/services/workspace_realtime/broadcaster_factory.py | 生产 broadcaster 单例 |
backend/services/workspace_realtime/workspace_presence_service.py | 在场态(Redis hash + snapshot) |
backend/services/workspace_realtime/canvas_event_handler.py | 画布操作落库编排 |
backend/services/workspace/workspace_canvas_event_service.py | 画布事件落库(含 output 布局更新 + 行锁) |
相关文档
docs/协作模式/2026-06-29-协作模式实现计划-后端.md:上游实现计划。docs/协作模式/2026-06-30-协作模式后端重构与实现设计.md:重构设计(God Class 拆分、Observer 挂接、Broadcaster 抽象的决策依据)。docs/协作模式/2026-06-29-协作模式实现计划-前端.md:前端协作侧约定。