Skip to content

协作模式后端实现总结:Agent 产物同步与实时协同 ​

本文是对协作模式后端已落地代码的总结解读,聚焦两件事:

  1. Agent 产物如何同步为画布产物(WorkspaceOutput)。
  2. 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 业务耦合。

采用两层抽象:

  1. 依赖倒置(Observer 协议):image/video/omni/chat 四个 stream 服务(加上 voice_service、action_mimic_service)在构造时注入一个 WorkspaceOutputObserver,默认是 NoopWorkspaceOutputObserver。stream 服务只认协议,不知道 workspace 存在与否。
  2. 策略族(产物同步来源):真正落库的逻辑封装成 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_startedstream 开始生成、已拿到 assistant_message_id建 pending 占位产物,广播 output:add
on_generation_completedstream 生成结束、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.py182 / 411started / completed
backend/services/core/video_stream_service.py643 / 660completed / video_task_submitted
backend/services/core/omni_stream_service.py205 / 393 / 410started / completed / video_task_submitted
backend/services/core/chat_stream_service.py406completed
backend/services/voice_service.py349 / 503completed / audio_task_submitted
backend/services/action_mimic_service.py215 / 351video_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 channel workspace:realtime:{workspace_id},收到消息后分发给本实例注册的本地 handler。

关键工程细节(backend/services/workspace_realtime/redis_broadcaster.py):

  1. publish 与 pubsub 用不同连接池。pubsub() 会独占一条连接直到关闭,若与 publish 共享 max_connections=50 的池,活跃 workspace 数达上限后会耗尽池导致 publish 失败。所以 pubsub 用独立池(_make_pubsub_redis,max_connections=64)。
  2. 每个 workspace 一个订阅协程,首条订阅时创建,最后一个 handler 退订时取消。
  3. 重连退避:_pubsub_loop 外层 while True,Redis 抖动/连接重置后 sleep(_RECONNECT_BACKOFF_SECONDS) 重订阅,保证跨实例广播在该实例持续可用(不"失明");仅 unsubscribe 触发的 CancelledError 才终止。
  4. 严格模式不降级:生产组合根 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/实例/workspaceConnectionRegistry.add防刷连接耗尽内存与 Redis 订阅
echo 排除sender_session_id发起方不收到自己乐观更新的事件回环
pubsub 重连退避_pubsub_loopRedis 抖动后自动重订阅,不"失明"
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.pyObserver 协议 + 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.pyAI 产物落库核心(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.pyREST 手动产物(查询/增改/软删/权限)
backend/models/workspace_output.py画布产物模型

实时协同子系统 ​

文件职责
backend/api/realtime/workspace_realtime.pyWS 端点薄编排(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.pyRedis 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:前端协作侧约定。