第2讲:OT 算法基础
2026/8/19 20:28:35 网站建设 项目流程

上一讲我们构建了文档引擎的核心——Piece Table 文本模型、操作历史、光标管理。但这些都是单用户视角。真正的协同编辑需要解决一个根本问题:

两个人同时编辑同一位置,如何保证最终结果一致?

这就是 Operational Transformation(OT 算法)要解决的问题。


一、OT 算法概述

1.1 问题场景

用户A和用户B同时编辑文档 "Hello" 用户A: 在位置5插入 " World" → 期望: "Hello World" 用户B: 在位置5插入 " !!!" → 期望: "Hello !!!" 如果简单按顺序应用: A先执行: "Hello" → "Hello World" B再执行: 在位置5插入 "!!!" → "Hello!!! World" ❌ 错误! 正确结果应该是: "Hello World!!!" 或 "Hello !!! World"

1.2 OT 的核心思想

OT 算法的本质: 当两个操作并发发生时,不直接应用对方的操作, 而是先将对方的操作"变换"(Transform)到自己的上下文上, 然后再应用。 操作变换公式: transform(op_a, op_b) → (op_a', op_b') 其中 op_a' 是在 op_b 之后执行时等价于 op_a 的操作 op_b' 是在 op_a 之后执行时等价于 op_b 的操作 性质:无论谁先执行,最终结果一致!

1.3 OT 的类型

类型

说明

代表系统

COT

客户端-服务端 OT

Google Wave

Jupiter

双向同步 OT

Google Docs 早期

GOTO

分组 OT

通用模型

Tombstone

墓碑变换

ShareJS

我们采用Jupiter 模型,它简洁且适合 WebSocket 双向通信。


二、操作定义与代数

2.1 操作的形式化

# core/ot/operation.py """ OT 操作的形式化定义 """ from __future__ import annotations from dataclasses import dataclass, field from enum import Enum from typing import Optional, Tuple import time import uuid class OpType(Enum): """操作类型""" INSERT = "insert" DELETE = "delete" NOOP = "noop" # 空操作 @dataclass class OTOperation: """ OT 操作 比基础 Operation 多了 site_id 和 sequence 信息, 用于分布式环境下的操作排序和追踪。 """ # 操作信息 op_type: OpType position: int text: str = "" deleted_text: str = "" # 分布式标识 site_id: str = "" # 产生此操作的站点 sequence: int = 0 # 站点内的操作序号 timestamp: float = 0.0 # 操作时间戳 # 操作依赖 depends_on: Optional[str] = None # 依赖的上一个操作 ID def __post_init__(self): if not self.timestamp: self.timestamp = time.time() if not self.site_id: self.site_id = f"site-{uuid.uuid4().hex[:8]}" @property def id(self) -> str: """操作唯一标识""" return f"{self.site_id}:{self.sequence}" def clone(self) -> OTOperation: """克隆操作""" return OTOperation( op_type=self.op_type, position=self.position, text=self.text, deleted_text=self.deleted_text, site_id=self.site_id, sequence=self.sequence, timestamp=self.timestamp ) def shift_position(self, delta: int): """偏移操作位置(用于变换)""" self.position += delta if self.position < 0: self.position = 0 def to_dict(self) -> dict: """序列化""" return { 'op_type': self.op_type.value, 'position': self.position, 'text': self.text, 'deleted_text': self.deleted_text, 'site_id': self.site_id, 'sequence': self.sequence, 'timestamp': self.timestamp, 'id': self.id } @classmethod def from_dict(cls, data: dict) -> OTOperation: """反序列化""" return cls( op_type=OpType(data['op_type']), position=data['position'], text=data.get('text', ''), deleted_text=data.get('deleted_text', ''), site_id=data.get('site_id', ''), sequence=data.get('sequence', 0), timestamp=data.get('timestamp', 0.0) ) def __repr__(self) -> str: if self.op_type == OpType.INSERT: return f"INS(+'{self.text}' @{self.position})[{self.id}]" elif self.op_type == OpType.DELETE: return f"DEL(-'{self.deleted_text}' @{self.position})[{self.id}]" return f"NOOP[{self.id}]"

三、变换函数实现

3.1 核心变换逻辑

# core/ot/transform.py """ OT 变换函数 实现操作之间的变换,保证并发操作的一致性。 """ from .operation import OTOperation, OpType import logging logger = logging.getLogger(__name__) def transform(op1: OTOperation, op2: OTOperation) -> tuple: """ 操作变换函数 将 op1 和 op2 变换为可以在对方之后执行的形式。 Args: op1: 第一个操作(假设已执行) op2: 第二个操作(待变换) Returns: (op1', op2') 变换后的操作对 其中 op1' 等价于在 op2 之后执行的 op1 op2' 等价于在 op1 之后执行的 op2 """ # 如果是 NOOP,直接返回 if op1.op_type == OpType.NOOP: return op1.clone(), op2.clone() if op2.op_type == OpType.NOOP: return op1.clone(), op2.clone() # 根据操作组合选择变换规则 if op1.op_type == OpType.INSERT and op2.op_type == OpType.INSERT: return _transform_insert_insert(op1, op2) elif op1.op_type == OpType.INSERT and op2.op_type == OpType.DELETE: return _transform_insert_delete(op1, op2) elif op1.op_type == OpType.DELETE and op2.op_type == OpType.INSERT: return _transform_delete_insert(op1, op2) elif op1.op_type == OpType.DELETE and op2.op_type == OpType.DELETE: return _transform_delete_delete(op1, op2) return op1.clone(), op2.clone() def _transform_insert_insert(op1: OTOperation, op2: OTOperation) -> tuple: """ INSERT vs INSERT 变换 两人在同一位置插入不同文本: - 如果 op1 的位置 < op2 的位置,op2 的位置后移 op1 的长度 - 如果 op1 的位置 > op2 的位置,op1 的位置后移 op2 的长度 - 如果位置相同,约定 site_id 小的先执行 """ op1_prime = op1.clone() op2_prime = op2.clone() if op1.position < op2.position: # op1 在 op2 前面插入,op2 的位置需要后移 op2_prime.shift_position(len(op1.text)) elif op1.position > op2.position: # op2 在 op1 前面插入,op1 的位置需要后移 op1_prime.shift_position(len(op2.text)) else: # 同一位置:按 site_id 排序 if op1.site_id < op2.site_id: op2_prime.shift_position(len(op1.text)) else: op1_prime.shift_position(len(op2.text)) logger.debug(f"XFORM I/I: {op1} + {op2} → ({op1_prime}, {op2_prime})") return op1_prime, op2_prime def _transform_insert_delete(op1: OTOperation, op2: OTOperation) -> tuple: """ INSERT vs DELETE 变换 插入操作遇到删除操作: - 如果插入位置在删除范围内,插入位置调整到删除起始位置 - 如果插入位置在删除范围之后,插入位置前移删除长度 """ op1_prime = op1.clone() op2_prime = op2.clone() delete_start = op2.position delete_end = op2.position + len(op2.deleted_text) if op1.position <= delete_start: # 插入在删除之前,不受影响 pass elif op1.position < delete_end: # 插入在删除范围内 → 移动到删除起始位置 op1_prime.position = delete_start else: # 插入在删除之后 → 位置前移 op1_prime.shift_position(-len(op2.deleted_text)) logger.debug(f"XFORM I/D: {op1} + {op2} → ({op1_prime}, {op2_prime})") return op1_prime, op2_prime def _transform_delete_insert(op1: OTOperation, op2: OTOperation) -> tuple: """ DELETE vs INSERT 变换 删除操作遇到插入操作: - 如果删除范围在插入之后,删除位置后移插入长度 - 如果删除范围包含插入位置,删除范围扩大 """ op1_prime = op1.clone() op2_prime = op2.clone() delete_start = op1.position delete_end = op1.position + len(op1.deleted_text) insert_len = len(op2.text) if op2.position <= delete_start: # 插入在删除之前 → 删除位置后移 op1_prime.shift_position(insert_len) elif op2.position < delete_end: # 插入在删除范围内 → 删除范围扩大 op1_prime.deleted_text = ( op1.deleted_text[:op2.position - delete_start] + op2.text + op1.deleted_text[op2.position - delete_start:] ) else: # 插入在删除之后,不受影响 pass logger.debug(f"XFORM D/I: {op1} + {op2} → ({op1_prime}, {op2_prime})") return op1_prime, op2_prime def _transform_delete_delete(op1: OTOperation, op2: OTOperation) -> tuple: """ DELETE vs DELETE 变换 两人删除同一区域: - 重叠部分只删除一次 - 非重叠部分各自调整位置 """ op1_prime = op1.clone() op2_prime = op2.clone() d1_start = op1.position d1_end = op1.position + len(op1.deleted_text) d2_start = op2.position d2_end = op2.position + len(op2.deleted_text) # 情况 1: op1 完全在 op2 之前 if d1_end <= d2_start: op2_prime.shift_position(-len(op1.deleted_text)) # 情况 2: op1 完全在 op2 之后 elif d1_start >= d2_end: op1_prime.shift_position(-len(op2.deleted_text)) # 情况 3: 有重叠 else: overlap_start = max(d1_start, d2_start) overlap_end = min(d1_end, d2_end) overlap_len = overlap_end - overlap_start # 调整 op1 的删除范围 if d1_start < d2_start: # op1 延伸到 op2 的起始 op1_prime.deleted_text = op1.deleted_text[:d2_start - d1_start] elif d1_start > d2_start: # op1 从 op2 内部开始 op1_prime.position = d2_start op1_prime.deleted_text = op1.deleted_text[d1_start - d2_start:] # 调整 op2 的删除范围 if d2_start < d1_start: op2_prime.deleted_text = op2.deleted_text[:d1_start - d2_start] elif d2_start > d1_start: op2_prime.position = d1_start op2_prime.deleted_text = op2.deleted_text[d2_start - d1_start:] logger.debug(f"XFORM D/D: {op1} + {op2} → ({op1_prime}, {op2_prime})") return op1_prime, op2_prime def compose(op1: OTOperation, op2: OTOperation) -> OTOperation: """ 操作组合 将两个连续的操作合并为一个等价操作。 用于批量同步时减少操作数量。 """ if op1.op_type == OpType.NOOP: return op2.clone() if op2.op_type == OpType.NOOP: return op1.clone() # INSERT + INSERT if op1.op_type == OpType.INSERT and op2.op_type == OpType.INSERT: if op2.position >= op1.position + len(op1.text): # op2 在 op1 后面 result = OTOperation( OpType.INSERT, op1.position, text=op1.text + op2.text, site_id=op1.site_id, sequence=op1.sequence ) return result # DELETE + DELETE if op1.op_type == OpType.DELETE and op2.op_type == OpType.DELETE: if op2.position == op1.position: # 连续删除 result = OTOperation( OpType.DELETE, op1.position, deleted_text=op1.deleted_text + op2.deleted_text, site_id=op1.site_id, sequence=op1.sequence ) return result # 无法组合,返回 op2 return op2.clone() def invert(op: OTOperation) -> OTOperation: """ 操作求逆 生成操作的逆操作,用于撤销。 """ if op.op_type == OpType.INSERT: return OTOperation( OpType.DELETE, op.position, deleted_text=op.text, site_id=op.site_id, sequence=op.sequence ) elif op.op_type == OpType.DELETE: return OTOperation( OpType.INSERT, op.position, text=op.deleted_text, site_id=op.site_id, sequence=op.sequence ) return OTOperation(OpType.NOOP, 0)

四、Jupiter 同步模型

4.1 模型架构

Jupiter 同步模型: 客户端 服务端 │ │ │─── op_A (seq=1) ──────▶│ (发送操作) │ │ │◀── ack (seq=1) ───────│ (确认) │ │ │─── op_B (seq=2) ──────▶│ │ │ │◀── op_C (from other) ──│ (接收远程操作) │ │ │─── transform(op_C) ────│ (变换后应用) │ │ 关键数据结构: - 客户端维护一个"未确认操作"列表(pending) - 服务端维护每个客户端的"已收到操作"列表 - 新操作到达时,与列表中所有操作进行变换

4.2 客户端同步引擎

# core/ot/sync_engine.py """ Jupiter 同步引擎 实现客户端和服务端的操作同步逻辑。 """ from __future__ import annotations from typing import List, Optional, Callable from .operation import OTOperation, OpType from .transform import transform import logging import threading import time logger = logging.getLogger(__name__) class SyncEngine: """ 同步引擎(Jupiter 模型) 负责: 1. 管理未确认操作列表 2. 执行操作变换 3. 维护操作顺序 """ def __init__(self, site_id: str, on_apply: Optional[Callable] = None): """ Args: site_id: 站点标识 on_apply: 应用操作的回调函数 """ self.site_id = site_id self.on_apply = on_apply # 未确认操作列表(等待服务端确认) self.pending: List[OTOperation] = [] # 已确认操作计数 self.confirmed_ops = 0 # 本地操作序号 self.local_seq = 0 # 锁 self.lock = threading.Lock() # ---------- 本地操作 ---------- def generate_local_op(self, op_type: OpType, position: int, text: str = "", deleted_text: str = "") -> OTOperation: """ 生成本地操作 Args: op_type: 操作类型 position: 位置 text: 插入文本 deleted_text: 删除文本 Returns: 生成的 OT 操作 """ with self.lock: self.local_seq += 1 op = OTOperation( op_type=op_type, position=position, text=text, deleted_text=deleted_text, site_id=self.site_id, sequence=self.local_seq ) # 加入未确认列表 self.pending.append(op) logger.debug(f"Local op generated: {op}") return op def apply_local_op(self, op: OTOperation): """应用本地操作(立即生效)""" if self.on_apply: self.on_apply(op) # ---------- 远程操作处理 ---------- def receive_remote_op(self, remote_op: OTOperation) -> Optional[OTOperation]: """ 接收远程操作 将远程操作与未确认操作进行变换,得到可以在本地执行的形式。 Args: remote_op: 远程操作 Returns: 变换后的操作(可在本地执行) """ with self.lock: transformed_op = remote_op.clone() # 与所有未确认操作进行变换 for pending_op in self.pending: # 变换: (pending_op, transformed_op) → (_, transformed_op') _, transformed_op = transform(pending_op, transformed_op) logger.debug(f"Remote op transformed: {remote_op} → {transformed_op}") return transformed_op def acknowledge_op(self, seq: int): """ 确认操作 服务端确认收到某个操作后,从未确认列表中移除。 """ with self.lock: # 找到并移除已确认的操作 for i, op in enumerate(self.pending): if op.sequence == seq: self.pending = self.pending[i + 1:] self.confirmed_ops += 1 logger.debug(f"Op acknowledged: {op}") break def get_pending_count(self) -> int: """获取未确认操作数""" with self.lock: return len(self.pending) def get_status(self) -> dict: """获取同步状态""" with self.lock: return { 'site_id': self.site_id, 'local_seq': self.local_seq, 'confirmed_ops': self.confirmed_ops, 'pending_count': len(self.pending) } class ServerSyncEngine: """ 服务端同步引擎 管理所有客户端的操作状态,执行全局变换。 """ def __init__(self, document_id: str): self.document_id = document_id # 所有已确认的操作(按顺序) self.history: List[OTOperation] = [] # 每个客户端的状态 self.client_states: dict = {} # site_id -> ClientState # 锁 self.lock = threading.Lock() def register_client(self, site_id: str) -> ClientState: """注册客户端""" state = ClientState(site_id) self.client_states[site_id] = state return state def receive_op(self, site_id: str, op: OTOperation) -> List[OTOperation]: """ 接收客户端操作 1. 与服务端历史进行变换 2. 广播给其他客户端 Args: site_id: 发送操作的客户端 op: 操作 Returns: 需要广播给其他客户端的操作列表 """ with self.lock: client_state = self.client_states.get(site_id) if not client_state: logger.warning(f"Unknown client: {site_id}") return [] # 与服务端历史进行变换 transformed_op = op.clone() for hist_op in self.history: _, transformed_op = transform(hist_op, transformed_op) # 添加到历史 self.history.append(transformed_op) # 更新客户端状态 client_state.last_seq = op.sequence # 生成需要广播给其他客户端的操作 broadcast_ops = [] for other_site, other_state in self.client_states.items(): if other_site != site_id: # 为其他客户端变换操作 other_op = transformed_op.clone() for hist_op in self.history[:-1]: # 排除自身 _, other_op = transform(hist_op, other_op) broadcast_ops.append(other_op) logger.debug(f"Server received {op}, broadcasting {len(broadcast_ops)} ops") return broadcast_ops def get_history_since(self, site_id: str, since_seq: int) -> List[OTOperation]: """ 获取某个客户端缺失的操作 用于客户端重连后同步。 """ with self.lock: # 找到该客户端最后一个确认的操作 client_state = self.client_states.get(site_id) if not client_state: return list(self.history) # 返回历史中该客户端之后的操作 missing_ops = [] for op in self.history: if op.site_id != site_id or op.sequence > client_state.last_seq: missing_ops.append(op) return missing_ops @dataclass class ClientState: """客户端状态""" site_id: str last_seq: int = 0 connected_at: float = 0.0 def __post_init__(self): if not self.connected_at: self.connected_at = time.time()

五、OT 集成到文档引擎

5.1 OTDocument 类

# core/ot/ot_document.py """ OT 增强的文档引擎 将 OT 同步集成到 Document 中。 """ from typing import List, Optional, Callable from ..document import Document from .operation import OTOperation, OpType from .sync_engine import SyncEngine from .transform import transform import logging logger = logging.getLogger(__name__) class OTDocument(Document): """ OT 文档引擎 继承自 Document,增加了 OT 同步能力。 """ def __init__(self, document_id: str, site_id: str, initial_text: str = ""): super().__init__(document_id, initial_text) self.site_id = site_id self.sync_engine = SyncEngine( site_id=site_id, on_apply=self._apply_op ) # 远程操作回调 self.on_remote_op: Optional[Callable] = None # ---------- 本地操作(OT 增强) ---------- def insert(self, position: int, text: str) -> OTOperation: """插入文本(生成 OT 操作)""" # 1. 生成 OT 操作 ot_op = self.sync_engine.generate_local_op( OpType.INSERT, position, text=text ) # 2. 应用本地操作 self.sync_engine.apply_local_op(ot_op) # 3. 更新文档内容 super().insert(position, text) return ot_op def delete(self, position: int, length: int) -> OTOperation: """删除文本(生成 OT 操作)""" # 获取被删除的文本 deleted_text = self.get_text_range(position, position + length) # 1. 生成 OT 操作 ot_op = self.sync_engine.generate_local_op( OpType.DELETE, position, deleted_text=deleted_text ) # 2. 应用本地操作 self.sync_engine.apply_local_op(ot_op) # 3. 更新文档内容 super().delete(position, length) return ot_op # ---------- 远程操作 ---------- def receive_remote_op(self, remote_op: OTOperation): """ 接收并应用远程操作 Args: remote_op: 来自其他站点的操作 """ # 1. 变换远程操作 transformed_op = self.sync_engine.receive_remote_op(remote_op) # 2. 应用到文档 self._apply_op(transformed_op) # 3. 通知回调 if self.on_remote_op: self.on_remote_op(transformed_op) def acknowledge_op(self, seq: int): """确认操作""" self.sync_engine.acknowledge_op(seq) def _apply_op(self, op: OTOperation): """应用操作到文档(不生成新的 OT 操作)""" if op.op_type == OpType.INSERT: super().insert(op.position, op.text) elif op.op_type == OpType.DELETE: super().delete(op.position, len(op.deleted_text)) def get_sync_status(self) -> dict: """获取同步状态""" return self.sync_engine.get_status()

六、测试

6.1 OT 变换测试

# tests/test_ot.py import pytest from core.ot.operation import OTOperation, OpType from core.ot.transform import ( transform, compose, invert, _transform_insert_insert, _transform_insert_delete, _transform_delete_insert, _transform_delete_delete ) class TestOTTransform: """OT 变换函数测试""" def test_insert_insert_same_position(self): """同一位置插入""" op1 = OTOperation(OpType.INSERT, 5, text="AB", site_id="A") op2 = OTOperation(OpType.INSERT, 5, text="CD", site_id="B") op1_prime, op2_prime = transform(op1, op2) # A 先执行,B 后移 assert op1_prime.position == 5 assert op2_prime.position == 7 # 5 + 2 def test_insert_insert_different_position(self): """不同位置插入""" op1 = OTOperation(OpType.INSERT, 3, text="XX", site_id="A") op2 = OTOperation(OpType.INSERT, 7, text="YY", site_id="B") op1_prime, op2_prime = transform(op1, op2) # op1 在 op2 前面,op2 后移 assert op1_prime.position == 3 assert op2_prime.position == 9 # 7 + 2 def test_insert_vs_delete_before(self): """插入在删除之前""" op1 = OTOperation(OpType.INSERT, 3, text="XX") op2 = OTOperation(OpType.DELETE, 5, deleted_text="Hello") op1_prime, op2_prime = transform(op1, op2) # 插入在删除之前,不受影响 assert op1_prime.position == 3 assert op2_prime.position == 5 def test_insert_vs_delete_inside(self): """插入在删除范围内""" op1 = OTOperation(OpType.INSERT, 6, text="XX") op2 = OTOperation(OpType.DELETE, 5, deleted_text="Hello") op1_prime, op2_prime = transform(op1, op2) # 插入在删除范围内,移到删除起始 assert op1_prime.position == 5 def test_delete_vs_delete_overlap(self): """重叠删除""" op1 = OTOperation(OpType.DELETE, 3, deleted_text="ABCDEF") op2 = OTOperation(OpType.DELETE, 5, deleted_text="DEFGHI") op1_prime, op2_prime = transform(op1, op2) # 重叠部分只删除一次 assert len(op1_prime.deleted_text) <= len(op1.deleted_text) class TestOTCompose: """操作组合测试""" def test_compose_inserts(self): """组合连续插入""" op1 = OTOperation(OpType.INSERT, 0, text="Hello") op2 = OTOperation(OpType.INSERT, 5, text=" World") result = compose(op1, op2) assert result.text == "Hello World" assert result.position == 0 class TestOTInvert: """操作求逆测试""" def test_invert_insert(self): """插入的逆操作是删除""" op = OTOperation(OpType.INSERT, 5, text="Hello") inv = invert(op) assert inv.op_type == OpType.DELETE assert inv.position == 5 assert inv.deleted_text == "Hello" def test_invert_delete(self): """删除的逆操作是插入""" op = OTOperation(OpType.DELETE, 3, deleted_text="World") inv = invert(op) assert inv.op_type == OpType.INSERT assert inv.position == 3 assert inv.text == "World" class TestOTDocument: """OT 文档测试""" def test_two_users_sync(self): """双用户同步测试""" from core.ot.ot_document import OTDocument # 用户 A 和 B 共享同一个文档 doc_a = OTDocument("doc-1", "site-A", "Hello") doc_b = OTDocument("doc-1", "site-B", "Hello") # 用户 A 插入 op_a = doc_a.insert(5, " World") # 用户 B 插入(并发) op_b = doc_b.insert(5, "!!!") # 互相接收对方操作 doc_a.receive_remote_op(op_b) doc_b.receive_remote_op(op_a) # 最终结果应该一致 assert doc_a.get_text() == doc_b.get_text() print(f"Doc A: '{doc_a.get_text()}'") print(f"Doc B: '{doc_b.get_text()}'") if __name__ == "__main__": pytest.main([__file__, "-v"])

七、可视化演示

7.1 OT 工作流演示

# examples/ot_demo.py """ OT 算法工作流演示 """ from core.ot.operation import OTOperation, OpType from core.ot.transform import transform from core.ot.sync_engine import SyncEngine, ServerSyncEngine from core.ot.ot_document import OTDocument def demonstrate_transform_rules(): """演示变换规则""" print("=" * 65) print("🔄 OT 变换规则演示") print("=" * 65) # 场景 1: 同一位置插入 print("\n📌 场景 1: 同一位置插入") print("-" * 40) doc = "Hello" op_a = OTOperation(OpType.INSERT, 5, text=" World", site_id="A", sequence=1) op_b = OTOperation(OpType.INSERT, 5, text="!!!", site_id="B", sequence=1) print(f"初始文档: '{doc}'") print(f"用户 A: {op_a}") print(f"用户 B: {op_b}") op_a_prime, op_b_prime = transform(op_a, op_b) print(f"\n变换后:") print(f" A' = {op_a_prime}") print(f" B' = {op_b_prime}") # 模拟执行 result_a = "Hello" + " World" + "!!!" # A 先插,B 后移 result_b = "Hello" + "!!!" + " World" # B 先插,A 后移 print(f"\n如果 A 先执行: '{result_a}'") print(f"如果 B 先执行: '{result_b}'") # 场景 2: 插入 vs 删除 print("\n📌 场景 2: 插入 vs 删除") print("-" * 40) doc = "Hello World" op_a = OTOperation(OpType.INSERT, 6, text="Beautiful ", site_id="A", sequence=1) op_b = OTOperation(OpType.DELETE, 5, deleted_text=" World", site_id="B", sequence=1) print(f"初始文档: '{doc}'") print(f"用户 A (插入): {op_a}") print(f"用户 B (删除): {op_b}") op_a_prime, op_b_prime = transform(op_a, op_b) print(f"\n变换后:") print(f" A' = {op_a_prime}") print(f" B' = {op_b_prime}") def simulate_collaboration(): """模拟协同编辑过程""" print("\n" + "=" * 65) print("👥 协同编辑模拟") print("=" * 65) # 创建两个用户和服务器 doc_a = OTDocument("demo", "Alice", "Hello") doc_b = OTDocument("demo", "Bob", "Hello") server = ServerSyncEngine("demo") server.register_client("Alice") server.register_client("Bob") print(f"\n初始文档: '{doc_a.get_text()}'") print(f"版本: {doc_a.version}") # Alice 打字 print("\n--- Alice 输入 ' World' ---") alice_op = doc_a.insert(5, " World") print(f"Alice 本地: '{doc_a.get_text()}'") # 发送到服务器 broadcast = server.receive_op("Alice", alice_op) doc_a.acknowledge_op(alice_op.sequence) # Bob 同时打字 print("\n--- Bob 同时输入 '!!!' ---") bob_op = doc_b.insert(5, "!!!") print(f"Bob 本地: '{doc_b.get_text()}'") # 发送到服务器 broadcast = server.receive_op("Bob", bob_op) doc_b.acknowledge_op(bob_op.sequence) # 互相接收 print("\n--- 同步 ---") doc_a.receive_remote_op(bob_op) doc_b.receive_remote_op(alice_op) print(f"\n最终 Alice: '{doc_a.get_text()}'") print(f"最终 Bob: '{doc_b.get_text()}'") print(f"一致: {doc_a.get_text() == doc_b.get_text()}") if __name__ == "__main__": demonstrate_transform_rules() simulate_collaboration()

八、总结

8.1 本讲成果

组件

文件

功能

OTOperation

core/ot/operation.py

分布式操作定义

Transform

core/ot/transform.py

四种变换规则

SyncEngine

core/ot/sync_engine.py

Jupiter 同步模型

ServerSyncEngine

core/ot/sync_engine.py

服务端同步引擎

OTDocument

core/ot/ot_document.py

OT 增强文档引擎

8.2 核心知识点

  • OT 的本质:通过操作变换保证并发操作的一致性

  • 四种变换规则:I/I、I/D、D/I、D/D 覆盖所有操作组合

  • Jupiter 模型:客户端维护 pending 列表,服务端维护全局历史

  • 最终一致性:无论操作顺序如何,所有站点最终结果一致

8.3 下一讲预告

第3讲:WebSocket 实时通信

我们将实现真正的网络通信层:

  • WebSocket 服务器与客户端

  • 操作消息的序列化与广播

  • 连接管理与心跳保活

  • 断线重连与状态恢复

准备好了吗?让我们在第3讲再见!


需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询