Спроектировать систему реального времени для одновременного редактирования документов несколькими пользователями (как Google Docs). Ключевой challenge -- разрешение конфликтов при одновременных изменениях.
Требования
Требование
Значение
Concurrent users
До 50 на один документ
Latency
< 100ms для локальных операций
Sync delay
< 500ms до отображения чужих правок
Consistency
Eventual consistency, convergence
Offline
Поддержка оффлайн-редактирования с синхронизацией
History
Полная история изменений с undo/redo
Подходы к синхронизации
OT (Operational Transformation)
OT -- алгоритм, который трансформирует операции относительно друг друга, чтобы результат применения был одинаковым независимо от порядка.
User A (position 2): Insert "X" → "abXcde"
User B (position 4): Insert "Y" → "abcdYe"
Без OT (naive): конфликт — позиции сдвигаются!
С OT:
1. Server получает Op(A) = Insert("X", pos=2)
2. Server получает Op(B) = Insert("Y", pos=4)
3. Transform: Op(B) после Op(A) → Insert("Y", pos=5) // +1 потому что X сдвинул
4. Result: "abXcdYe" — одинаковый у обоих!
CRDT (Conflict-free Replicated Data Type)
CRDT -- структуры данных, которые гарантируют конвергенцию без центрального сервера. Каждая реплика может применять операции независимо.
User A: вставляет "X" после символа с ID "a2"
User B: вставляет "Y" после символа с ID "a4"
CRDT: каждый символ имеет уникальный ID
Результат: операции коммутативны, порядок не важен
Все реплики сходятся к одному состоянию
<?php
declare(strict_types=1);
namespace App\Collaboration;
/**
* Document synchronization service.
* Receives operations from clients and broadcasts to others.
*/
final class DocumentSyncService
{
/** @var array<string, array<string>> Active connections per document */
private array $documentSessions = [];
/** @var array<string, int> Server version per document */
private array $documentVersions = [];
public function __construct(
private readonly OperationStore $operationStore,
private readonly DocumentStore $documentStore,
) {}
/**
* Process an operation from a client.
*
* @param string $documentId
* @param string $clientId
* @param Operation $operation
* @return SyncResult Operations to broadcast to other clients
*/
public function processOperation(
string $documentId,
string $clientId,
Operation $operation,
): SyncResult {
$currentVersion = $this->documentVersions[$documentId] ?? 0;
// If client is behind, transform operation against missed ops
if ($operation->baseVersion < $currentVersion) {
$missedOps = $this->operationStore->getOperationsSince(
$documentId,
$operation->baseVersion,
);
// Transform the incoming operation against each missed operation
foreach ($missedOps as $missedOp) {
if ($missedOp->clientId === $clientId) {
continue; // Skip own operations
}
$operation = $this->transform($operation, $missedOp);
}
}
// Assign server version
$newVersion = $currentVersion + 1;
$this->documentVersions[$documentId] = $newVersion;
// Persist operation
$serverOp = new ServerOperation(
documentId: $documentId,
clientId: $clientId,
operation: $operation,
serverVersion: $newVersion,
timestamp: new \DateTimeImmutable(),
);
$this->operationStore->save($serverOp);
// Return acknowledgement for sender + broadcast for others
return new SyncResult(
acknowledgement: new Acknowledgement(
clientId: $clientId,
serverVersion: $newVersion,
),
broadcast: new BroadcastOperation(
operation: $operation,
serverVersion: $newVersion,
originClientId: $clientId,
),
recipients: $this->getOtherClients($documentId, $clientId),
);
}
/**
* Transform operation A against operation B (OT core).
* This is simplified — real OT has many edge cases.
*/
private function transform(Operation $incoming, ServerOperation $existing): Operation
{
// Both are inserts
if ($incoming->type === OperationType::Insert && $existing->operation->type === OperationType::Insert) {
if ($existing->operation->position <= $incoming->position) {
// Existing insert shifts incoming position right
return new Operation(
type: $incoming->type,
position: $incoming->position + strlen($existing->operation->content),
content: $incoming->content,
baseVersion: $incoming->baseVersion,
);
}
}
// Existing delete before incoming insert
if ($existing->operation->type === OperationType::Delete && $incoming->type === OperationType::Insert) {
if ($existing->operation->position < $incoming->position) {
return new Operation(
type: $incoming->type,
position: $incoming->position - $existing->operation->length,
content: $incoming->content,
baseVersion: $incoming->baseVersion,
);
}
}
return $incoming; // No transformation needed
}
/**
* Create a document snapshot for new clients joining.
*/
public function getDocumentSnapshot(string $documentId): DocumentSnapshot
{
$document = $this->documentStore->get($documentId);
$version = $this->documentVersions[$documentId] ?? 0;
return new DocumentSnapshot(
documentId: $documentId,
content: $document->getContent(),
version: $version,
cursors: $this->getActiveCursors($documentId),
);
}
private function getOtherClients(string $documentId, string $excludeClientId): array
{
$clients = $this->documentSessions[$documentId] ?? [];
return array_filter($clients, static fn(string $id) => $id !== $excludeClientId);
}
private function getActiveCursors(string $documentId): array
{
return []; // Cursor positions of active users
}
}
package collaboration
import (
"sync"
"time"
)
// DocumentSyncService handles document synchronization via OT.
type DocumentSyncService struct {
mu sync.Mutex
documentSessions map[string][]string // documentID -> clientIDs
documentVersions map[string]int // documentID -> version
operationStore OperationStore
documentStore DocumentStore
}
// NewDocumentSyncService creates a sync service.
func NewDocumentSyncService(opStore OperationStore, docStore DocumentStore) *DocumentSyncService {
return &DocumentSyncService{
documentSessions: make(map[string][]string),
documentVersions: make(map[string]int),
operationStore: opStore,
documentStore: docStore,
}
}
// ProcessOperation receives a client operation, transforms it if needed, and returns a sync result.
func (s *DocumentSyncService) ProcessOperation(documentID, clientID string, op Operation) (SyncResult, error) {
s.mu.Lock()
defer s.mu.Unlock()
currentVersion := s.documentVersions[documentID]
// Transform against missed operations
if op.BaseVersion < currentVersion {
missedOps, err := s.operationStore.GetOperationsSince(documentID, op.BaseVersion)
if err != nil {
return SyncResult{}, err
}
for _, missed := range missedOps {
if missed.ClientID == clientID {
continue
}
op = transform(op, missed.Operation)
}
}
newVersion := currentVersion + 1
s.documentVersions[documentID] = newVersion
serverOp := ServerOperation{
DocumentID: documentID,
ClientID: clientID,
Operation: op,
ServerVersion: newVersion,
Timestamp: time.Now(),
}
if err := s.operationStore.Save(serverOp); err != nil {
return SyncResult{}, err
}
return SyncResult{
Acknowledgement: Acknowledgement{ClientID: clientID, ServerVersion: newVersion},
Broadcast: BroadcastOperation{Operation: op, ServerVersion: newVersion, OriginClientID: clientID},
Recipients: s.getOtherClients(documentID, clientID),
}, nil
}
func transform(incoming Operation, existing Operation) Operation {
if incoming.Type == OpInsert && existing.Type == OpInsert {
if existing.Position <= incoming.Position {
incoming.Position += len(existing.Content)
}
}
if existing.Type == OpDelete && incoming.Type == OpInsert {
if existing.Position < incoming.Position {
incoming.Position -= existing.Length
}
}
return incoming
}
func (s *DocumentSyncService) getOtherClients(documentID, excludeID string) []string {
var others []string
for _, id := range s.documentSessions[documentID] {
if id != excludeID {
others = append(others, id)
}
}
return others
}
namespace App.Collaboration;
/// Document synchronization service.
/// Receives operations from clients and broadcasts them to others.
public sealed class DocumentSyncService
{
private readonly Dictionary<string, List<string>> _documentSessions = [];
private readonly Dictionary<string, int> _documentVersions = [];
private readonly SemaphoreSlim _lock = new(1, 1);
private readonly IOperationStore _operationStore;
private readonly IDocumentStore _documentStore;
public DocumentSyncService(IOperationStore operationStore, IDocumentStore documentStore)
{
_operationStore = operationStore;
_documentStore = documentStore;
}
// Process an operation from a client and return what to broadcast.
public async Task<SyncResult> ProcessOperationAsync(
string documentId,
string clientId,
Operation operation,
CancellationToken ct = default)
{
await _lock.WaitAsync(ct);
try
{
var currentVersion = _documentVersions.GetValueOrDefault(documentId);
// If the client is behind, transform the operation against missed ops
if (operation.BaseVersion < currentVersion)
{
var missedOps = await _operationStore.GetOperationsSinceAsync(
documentId,
operation.BaseVersion,
ct);
foreach (var missed in missedOps)
{
if (missed.ClientId == clientId)
{
continue; // Skip own operations
}
operation = Transform(operation, missed.Operation);
}
}
// Assign server version
var newVersion = currentVersion + 1;
_documentVersions[documentId] = newVersion;
// Persist operation
var serverOp = new ServerOperation(
documentId,
clientId,
operation,
newVersion,
DateTimeOffset.UtcNow);
await _operationStore.SaveAsync(serverOp, ct);
// Acknowledgement for sender + broadcast for everyone else
return new SyncResult(
new Acknowledgement(clientId, newVersion),
new BroadcastOperation(operation, newVersion, clientId),
GetOtherClients(documentId, clientId));
}
finally
{
_lock.Release();
}
}
// Create a document snapshot for new clients joining.
public async Task<DocumentSnapshot> GetDocumentSnapshotAsync(
string documentId,
CancellationToken ct = default)
{
var document = await _documentStore.GetAsync(documentId, ct);
var version = _documentVersions.GetValueOrDefault(documentId);
return document with { Version = version };
}
/// Transform the incoming operation against an existing one (OT core).
/// Simplified — real OT has many more edge cases.
private static Operation Transform(Operation incoming, Operation existing)
{
// Both are inserts: an earlier insert shifts the incoming position right
if (incoming.Type is OperationType.Insert
&& existing.Type is OperationType.Insert
&& existing.Position <= incoming.Position)
{
return incoming with { Position = incoming.Position + (existing.Content?.Length ?? 0) };
}
// Existing delete before incoming insert: shift the position left
if (existing.Type is OperationType.Delete
&& incoming.Type is OperationType.Insert
&& existing.Position < incoming.Position)
{
return incoming with { Position = incoming.Position - (existing.Length ?? 0) };
}
return incoming; // No transformation needed
}
private List<string> GetOtherClients(string documentId, string excludeClientId)
=> _documentSessions.TryGetValue(documentId, out var clients)
? clients.Where(id => id != excludeClientId).ToList()
: [];
}
import asyncio
from collections import defaultdict
from dataclasses import replace
from datetime import datetime, timezone
class DocumentSyncService:
"""Document synchronization service.
Receives operations from clients and broadcasts them to others.
"""
def __init__(
self,
operation_store: OperationStore,
document_store: DocumentStore,
) -> None:
self._operation_store = operation_store
self._document_store = document_store
self._document_sessions: dict[str, list[str]] = defaultdict(list)
self._document_versions: dict[str, int] = defaultdict(int)
# asyncio.Lock replaces Go's mutex — the event loop is single-threaded,
# but await points still interleave coroutines.
self._lock = asyncio.Lock()
async def process_operation(
self,
document_id: str,
client_id: str,
operation: Operation,
) -> SyncResult:
"""Process an operation from a client and return what to broadcast."""
async with self._lock:
current_version = self._document_versions[document_id]
# If the client is behind, transform the operation against missed ops
if operation.base_version < current_version:
missed_ops = await self._operation_store.get_operations_since(
document_id, operation.base_version
)
for missed in missed_ops:
if missed.client_id == client_id:
continue # Skip own operations
operation = self._transform(operation, missed.operation)
# Assign server version
new_version = current_version + 1
self._document_versions[document_id] = new_version
# Persist operation
server_op = ServerOperation(
document_id=document_id,
client_id=client_id,
operation=operation,
server_version=new_version,
timestamp=datetime.now(timezone.utc),
)
await self._operation_store.save(server_op)
# Acknowledgement for sender + broadcast for everyone else
return SyncResult(
acknowledgement=Acknowledgement(client_id, new_version),
broadcast=BroadcastOperation(operation, new_version, client_id),
recipients=self._other_clients(document_id, client_id),
)
async def get_document_snapshot(self, document_id: str) -> DocumentSnapshot:
"""Create a document snapshot for new clients joining."""
snapshot = await self._document_store.get(document_id)
return replace(snapshot, version=self._document_versions[document_id])
@staticmethod
def _transform(incoming: Operation, existing: Operation) -> Operation:
"""Transform the incoming operation against an existing one (OT core).
Simplified — real OT has many more edge cases.
"""
# Both are inserts: an earlier insert shifts the incoming position right
if (
incoming.type is OperationType.INSERT
and existing.type is OperationType.INSERT
and existing.position <= incoming.position
):
return replace(
incoming,
position=incoming.position + len(existing.content or ""),
)
# Existing delete before incoming insert: shift the position left
if (
existing.type is OperationType.DELETE
and incoming.type is OperationType.INSERT
and existing.position < incoming.position
):
return replace(incoming, position=incoming.position - (existing.length or 0))
return incoming # No transformation needed
def _other_clients(self, document_id: str, exclude_client_id: str) -> list[str]:
return [
client_id
for client_id in self._document_sessions[document_id]
if client_id != exclude_client_id
]
## Типы операций
<?php
declare(strict_types=1);
namespace App\Collaboration;
enum OperationType: string
{
case Insert = 'insert';
case Delete = 'delete';
case Retain = 'retain'; // Skip N characters
case Format = 'format'; // Bold, italic, etc.
}
final readonly class Operation
{
public function __construct(
public OperationType $type,
public int $position,
public ?string $content = null, // For insert
public ?int $length = null, // For delete/retain
public ?array $attributes = null, // For format
public int $baseVersion = 0,
) {}
}
final readonly class ServerOperation
{
public function __construct(
public string $documentId,
public string $clientId,
public Operation $operation,
public int $serverVersion,
public \DateTimeImmutable $timestamp,
) {}
}
final readonly class Acknowledgement
{
public function __construct(
public string $clientId,
public int $serverVersion,
) {}
}
final readonly class BroadcastOperation
{
public function __construct(
public Operation $operation,
public int $serverVersion,
public string $originClientId,
) {}
}
final readonly class SyncResult
{
public function __construct(
public Acknowledgement $acknowledgement,
public BroadcastOperation $broadcast,
/** @var array<string> */
public array $recipients,
) {}
}
final readonly class DocumentSnapshot
{
public function __construct(
public string $documentId,
public string $content,
public int $version,
public array $cursors,
) {}
}
package collaboration
import "time"
// Operation types.
const (
OpInsert = "insert"
OpDelete = "delete"
OpRetain = "retain"
OpFormat = "format"
)
// Operation represents an edit operation.
type Operation struct {
Type string
Position int
Content string // For insert
Length int // For delete/retain
Attributes map[string]string // For format
BaseVersion int
}
// ServerOperation wraps an operation with server metadata.
type ServerOperation struct {
DocumentID string
ClientID string
Operation Operation
ServerVersion int
Timestamp time.Time
}
// Acknowledgement confirms operation receipt.
type Acknowledgement struct {
ClientID string
ServerVersion int
}
// BroadcastOperation is sent to other clients.
type BroadcastOperation struct {
Operation Operation
ServerVersion int
OriginClientID string
}
// SyncResult holds the result of processing an operation.
type SyncResult struct {
Acknowledgement Acknowledgement
Broadcast BroadcastOperation
Recipients []string
}
// DocumentSnapshot captures the document state at a point in time.
type DocumentSnapshot struct {
DocumentID string
Content string
Version int
Cursors map[string]CursorPosition
}
// OperationStore persists operations.
type OperationStore interface {
Save(op ServerOperation) error
GetOperationsSince(documentID string, version int) ([]ServerOperation, error)
}
// DocumentStore manages document snapshots.
type DocumentStore interface {
Get(documentID string) (*DocumentSnapshot, error)
}
namespace App.Collaboration;
public enum OperationType
{
Insert,
Delete,
Retain, // Skip N characters
Format, // Bold, italic, etc.
}
/// Represents a single edit operation.
public sealed record Operation(
OperationType Type,
int Position,
string? Content = null, // For insert
int? Length = null, // For delete/retain
IReadOnlyDictionary<string, string>? Attributes = null, // For format
int BaseVersion = 0);
/// Wraps an operation with server metadata.
public sealed record ServerOperation(
string DocumentId,
string ClientId,
Operation Operation,
int ServerVersion,
DateTimeOffset Timestamp);
/// Confirms operation receipt to the sender.
public sealed record Acknowledgement(string ClientId, int ServerVersion);
/// Sent to all other clients editing the document.
public sealed record BroadcastOperation(
Operation Operation,
int ServerVersion,
string OriginClientId);
/// Holds the result of processing an operation.
public sealed record SyncResult(
Acknowledgement Acknowledgement,
BroadcastOperation Broadcast,
IReadOnlyList<string> Recipients);
/// Captures the document state at a point in time.
public sealed record DocumentSnapshot(
string DocumentId,
string Content,
int Version,
IReadOnlyDictionary<string, CursorPosition> Cursors);
/// Persists operations.
public interface IOperationStore
{
Task SaveAsync(ServerOperation op, CancellationToken ct = default);
Task<IReadOnlyList<ServerOperation>> GetOperationsSinceAsync(
string documentId,
int version,
CancellationToken ct = default);
}
/// Manages document snapshots.
public interface IDocumentStore
{
Task<DocumentSnapshot> GetAsync(string documentId, CancellationToken ct = default);
}
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from typing import Protocol, Sequence
class OperationType(str, Enum):
INSERT = "insert"
DELETE = "delete"
RETAIN = "retain" # Skip N characters
FORMAT = "format" # Bold, italic, etc.
@dataclass(frozen=True, slots=True)
class Operation:
"""Represents a single edit operation."""
type: OperationType
position: int
content: str | None = None # For insert
length: int | None = None # For delete/retain
attributes: dict[str, str] | None = None # For format
base_version: int = 0
@dataclass(frozen=True, slots=True)
class ServerOperation:
"""Wraps an operation with server metadata."""
document_id: str
client_id: str
operation: Operation
server_version: int
timestamp: datetime
@dataclass(frozen=True, slots=True)
class Acknowledgement:
"""Confirms operation receipt to the sender."""
client_id: str
server_version: int
@dataclass(frozen=True, slots=True)
class BroadcastOperation:
"""Sent to all other clients editing the document."""
operation: Operation
server_version: int
origin_client_id: str
@dataclass(frozen=True, slots=True)
class SyncResult:
"""Holds the result of processing an operation."""
acknowledgement: Acknowledgement
broadcast: BroadcastOperation
recipients: list[str]
@dataclass(frozen=True, slots=True)
class DocumentSnapshot:
"""Captures the document state at a point in time."""
document_id: str
content: str
version: int
cursors: dict[str, "CursorPosition"] = field(default_factory=dict)
class OperationStore(Protocol):
# Python has no interfaces; typing.Protocol describes the required
# methods structurally, so any storage backend can satisfy it.
async def save(self, op: ServerOperation) -> None: ...
async def get_operations_since(
self,
document_id: str,
version: int,
) -> Sequence[ServerOperation]: ...
class DocumentStore(Protocol):
async def get(self, document_id: str) -> DocumentSnapshot: ...
## Presence (курсоры пользователей)
<?php
declare(strict_types=1);
namespace App\Collaboration;
final class PresenceService
{
/** @var array<string, array<string, CursorPosition>> */
private array $cursors = [];
/**
* Update cursor position for a user.
*/
public function updateCursor(
string $documentId,
string $userId,
CursorPosition $position,
): void {
$this->cursors[$documentId][$userId] = $position;
}
/**
* Get all active cursors for a document.
*
* @return array<string, CursorPosition>
*/
public function getCursors(string $documentId): array
{
return $this->cursors[$documentId] ?? [];
}
/**
* Remove cursor when user disconnects.
*/
public function removeCursor(string $documentId, string $userId): void
{
unset($this->cursors[$documentId][$userId]);
}
}
final readonly class CursorPosition
{
public function __construct(
public int $offset,
public ?int $selectionEnd = null,
public string $userName = '',
public string $color = '#FF0000',
) {}
}
package collaboration
import "sync"
// CursorPosition represents a user's cursor location.
type CursorPosition struct {
Offset int
SelectionEnd *int
UserName string
Color string
}
// PresenceService tracks active user cursors in documents.
type PresenceService struct {
mu sync.RWMutex
cursors map[string]map[string]CursorPosition // documentID -> userID -> position
}
// NewPresenceService creates a presence service.
func NewPresenceService() *PresenceService {
return &PresenceService{cursors: make(map[string]map[string]CursorPosition)}
}
// UpdateCursor sets the cursor position for a user.
func (s *PresenceService) UpdateCursor(documentID, userID string, pos CursorPosition) {
s.mu.Lock()
defer s.mu.Unlock()
if s.cursors[documentID] == nil {
s.cursors[documentID] = make(map[string]CursorPosition)
}
s.cursors[documentID][userID] = pos
}
// GetCursors returns all active cursors for a document.
func (s *PresenceService) GetCursors(documentID string) map[string]CursorPosition {
s.mu.RLock()
defer s.mu.RUnlock()
result := make(map[string]CursorPosition)
for k, v := range s.cursors[documentID] {
result[k] = v
}
return result
}
// RemoveCursor removes a user's cursor on disconnect.
func (s *PresenceService) RemoveCursor(documentID, userID string) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.cursors[documentID], userID)
}
using System.Collections.Concurrent;
namespace App.Collaboration;
/// Represents a user's cursor location in a document.
public sealed record CursorPosition(
int Offset,
int? SelectionEnd = null,
string UserName = "",
string Color = "#FF0000");
/// Tracks active user cursors per document.
public sealed class PresenceService
{
// ConcurrentDictionary replaces an explicit lock — cursor updates
// arrive from many WebSocket connections at once.
private readonly ConcurrentDictionary<string, ConcurrentDictionary<string, CursorPosition>> _cursors = new();
// Update the cursor position for a user.
public void UpdateCursor(string documentId, string userId, CursorPosition position)
=> _cursors.GetOrAdd(documentId, _ => new ConcurrentDictionary<string, CursorPosition>())[userId] = position;
// Get all active cursors for a document.
public IReadOnlyDictionary<string, CursorPosition> GetCursors(string documentId)
=> _cursors.TryGetValue(documentId, out var cursors)
? new Dictionary<string, CursorPosition>(cursors)
: new Dictionary<string, CursorPosition>();
// Remove a user's cursor when they disconnect.
public void RemoveCursor(string documentId, string userId)
{
if (_cursors.TryGetValue(documentId, out var cursors))
{
cursors.TryRemove(userId, out _);
}
}
}
from collections import defaultdict
from dataclasses import dataclass
@dataclass(frozen=True, slots=True)
class CursorPosition:
"""Represents a user's cursor location in a document."""
offset: int
selection_end: int | None = None
user_name: str = ""
color: str = "#FF0000"
class PresenceService:
"""Track active user cursors per document.
No lock is needed: dict mutations are atomic under the GIL and the
asyncio event loop never interleaves these synchronous methods.
"""
def __init__(self) -> None:
self._cursors: dict[str, dict[str, CursorPosition]] = defaultdict(dict)
def update_cursor(
self,
document_id: str,
user_id: str,
position: CursorPosition,
) -> None:
"""Update the cursor position for a user."""
self._cursors[document_id][user_id] = position
def get_cursors(self, document_id: str) -> dict[str, CursorPosition]:
"""Get all active cursors for a document."""
return dict(self._cursors.get(document_id, {}))
def remove_cursor(self, document_id: str, user_id: str) -> None:
"""Remove a user's cursor when they disconnect."""
self._cursors.get(document_id, {}).pop(user_id, None)
## Конфликтные сценарии
Сценарий
Решение
Два пользователя пишут в одно место
OT трансформирует позиции
Один удаляет текст, другой редактирует его
OT: delete wins, edit отбрасывается
Оффлайн правки + онлайн правки
При reconnect: трансформировать все оффлайн-операции
Конфликт форматирования
Last-write-wins для атрибутов
Одновременное удаление
Idempotent delete (повторное удаление = no-op)
Хранение и History
Snapshot + Operations model
Snapshot (version 0): "Hello"
Op v1: Insert("World", pos=5) → "HelloWorld"
Op v2: Insert(" ", pos=5) → "Hello World"
Op v3: Delete(pos=0, len=5) → " World"
Recovery: Apply snapshot + all ops since snapshot
Optimization: Periodic snapshot creation (every 100 ops)
<?php
declare(strict_types=1);
namespace App\Collaboration;
final readonly class DocumentHistory
{
public function __construct(
private OperationStore $operations,
private DocumentStore $documents,
) {}
/**
* Rebuild document state at a specific version.
*/
public function getAtVersion(string $documentId, int $version): string
{
// Find nearest snapshot before requested version
$snapshot = $this->documents->getSnapshotBefore($documentId, $version);
// Apply operations from snapshot to target version
$ops = $this->operations->getOperationsRange(
$documentId,
fromVersion: $snapshot->version,
toVersion: $version,
);
$content = $snapshot->content;
foreach ($ops as $op) {
$content = $this->applyOperation($content, $op->operation);
}
return $content;
}
private function applyOperation(string $content, Operation $op): string
{
return match ($op->type) {
OperationType::Insert => substr($content, 0, $op->position)
. $op->content
. substr($content, $op->position),
OperationType::Delete => substr($content, 0, $op->position)
. substr($content, $op->position + $op->length),
default => $content,
};
}
}
package collaboration
// DocumentHistory rebuilds document state from operations.
type DocumentHistory struct {
operations OperationStore
documents DocumentStore
}
// NewDocumentHistory creates a document history service.
func NewDocumentHistory(ops OperationStore, docs DocumentStore) *DocumentHistory {
return &DocumentHistory{operations: ops, documents: docs}
}
// GetAtVersion rebuilds the document at a specific version.
func (h *DocumentHistory) GetAtVersion(documentID string, version int) (string, error) {
snapshot, err := h.documents.Get(documentID)
if err != nil {
return "", err
}
ops, err := h.operations.GetOperationsSince(documentID, snapshot.Version)
if err != nil {
return "", err
}
content := snapshot.Content
for _, op := range ops {
if op.ServerVersion > version {
break
}
content = applyOperation(content, op.Operation)
}
return content, nil
}
func applyOperation(content string, op Operation) string {
runes := []rune(content)
switch op.Type {
case OpInsert:
if op.Position > len(runes) {
op.Position = len(runes)
}
result := make([]rune, 0, len(runes)+len([]rune(op.Content)))
result = append(result, runes[:op.Position]...)
result = append(result, []rune(op.Content)...)
result = append(result, runes[op.Position:]...)
return string(result)
case OpDelete:
end := op.Position + op.Length
if end > len(runes) {
end = len(runes)
}
result := make([]rune, 0, len(runes)-op.Length)
result = append(result, runes[:op.Position]...)
result = append(result, runes[end:]...)
return string(result)
default:
return content
}
}
namespace App.Collaboration;
/// Rebuilds document state from a snapshot plus its operation log.
public sealed class DocumentHistory
{
private readonly IOperationStore _operations;
private readonly IDocumentStore _documents;
public DocumentHistory(IOperationStore operations, IDocumentStore documents)
{
_operations = operations;
_documents = documents;
}
// Rebuild the document state at a specific version.
public async Task<string> GetAtVersionAsync(
string documentId,
int version,
CancellationToken ct = default)
{
// Find the nearest snapshot before the requested version
var snapshot = await _documents.GetAsync(documentId, ct);
// Apply operations from the snapshot up to the target version
var ops = await _operations.GetOperationsSinceAsync(documentId, snapshot.Version, ct);
var content = snapshot.Content;
foreach (var op in ops.Where(op => op.ServerVersion <= version))
{
content = ApplyOperation(content, op.Operation);
}
return content;
}
private static string ApplyOperation(string content, Operation op) => op.Type switch
{
OperationType.Insert => content
.Insert(Math.Min(op.Position, content.Length), op.Content ?? string.Empty),
OperationType.Delete => content
.Remove(
Math.Min(op.Position, content.Length),
Math.Min(op.Length ?? 0, Math.Max(0, content.Length - op.Position))),
// Retain and Format do not change the plain-text content
OperationType.Retain or OperationType.Format => content,
_ => throw new ArgumentOutOfRangeException(
nameof(op),
op.Type,
"Unknown operation type"),
};
}
class DocumentHistory:
"""Rebuild document state from a snapshot plus its operation log."""
def __init__(self, operations: OperationStore, documents: DocumentStore) -> None:
self._operations = operations
self._documents = documents
async def get_at_version(self, document_id: str, version: int) -> str:
"""Rebuild the document state at a specific version."""
# Find the nearest snapshot before the requested version
snapshot = await self._documents.get(document_id)
# Apply operations from the snapshot up to the target version
ops = await self._operations.get_operations_since(document_id, snapshot.version)
content = snapshot.content
for op in ops:
if op.server_version > version:
break
content = self._apply_operation(content, op.operation)
return content
@staticmethod
def _apply_operation(content: str, op: Operation) -> str:
position = min(op.position, len(content))
match op.type:
case OperationType.INSERT:
return content[:position] + (op.content or "") + content[position:]
case OperationType.DELETE:
return content[:position] + content[position + (op.length or 0) :]
case OperationType.RETAIN | OperationType.FORMAT:
# Neither changes the plain-text content
return content
case _:
raise ValueError(f"Unknown operation type: {op.type}")
## Масштабирование
Аспект
Решение
Много документов
Шардирование по document_id
Много пользователей на документ
Один сервер на документ (affinity)
Персистентность
Периодические snapshots + append-only operations
Failover
Rebuild из operations log
WebSocket масштабирование
Redis pub/sub для cross-server broadcast
Итоги
Концепция
Суть
OT
Трансформация операций для консистентности
CRDT
Структуры данных с автоматической конвергенцией
Cursor-based sync
Каждая операция привязана к позиции
Presence
Отображение курсоров других пользователей
Snapshot + Ops
Эффективное хранение истории
Conflict resolution
OT transform или CRDT merge
Главный вывод: Коллаборативное редактирование -- одна из сложнейших задач в distributed systems. Начинайте с OT (проще для понимания, нужен сервер), переходите к CRDT если нужен offline и P2P.