278 lines
9.5 KiB
JavaScript
278 lines
9.5 KiB
JavaScript
/*
|
|
* Copyright (C) 2026 Mateusz Gruszczynski @linuxiarz.pl
|
|
* Source-Available Code / Dual-Licensed.
|
|
*
|
|
* Free for non-commercial and evaluation use under terms of BSL/GPLv3.
|
|
* Commercial or production use requires a valid paid license.
|
|
* See LICENSE file in repository root for details.
|
|
*/
|
|
|
|
import {
|
|
applyOperationToDocument,
|
|
compareOperationKeys,
|
|
composeOperations,
|
|
documentAfterPending,
|
|
identityOperation,
|
|
normalizeOperation,
|
|
operationBaseLength,
|
|
operationFromEdit,
|
|
operationTargetLength,
|
|
transformOperations,
|
|
} from "@rustpad/collaboration";
|
|
import { parseAuthorship } from "@rustpad/authorship";
|
|
|
|
function hasTextEffect(operation) {
|
|
return normalizeOperation(operation).components.some(component => component.kind === "insert" || component.kind === "delete");
|
|
}
|
|
|
|
function mergeOwnerReplacements(previous = [], next = []) {
|
|
const replacements = new Map();
|
|
for (const item of [...previous, ...next]) {
|
|
const owner = String(item?.owner || "");
|
|
const replacement = String(item?.replacement || "");
|
|
if (owner && replacement) replacements.set(owner, { owner, replacement });
|
|
}
|
|
return [...replacements.values()];
|
|
}
|
|
|
|
function pendingEnvelope(clientId, updateId, operation, ownerReplacements = []) {
|
|
return {
|
|
clientId,
|
|
updateId,
|
|
operation: normalizeOperation(operation),
|
|
ownerReplacements: mergeOwnerReplacements([], ownerReplacements),
|
|
};
|
|
}
|
|
|
|
function serverEnvelope(message) {
|
|
return {
|
|
baseRevisionId: Number(message?.base_revision_id ?? message?.baseRevisionId),
|
|
revisionId: Number(message?.revision_id ?? message?.revisionId),
|
|
clientId: String(message?.client_id ?? message?.clientId ?? ""),
|
|
updateId: Number(message?.update_id ?? message?.updateId ?? 0),
|
|
operation: normalizeOperation(message?.operation),
|
|
ownerReplacements: mergeOwnerReplacements([], message?.owner_replacements ?? message?.ownerReplacements ?? []),
|
|
};
|
|
}
|
|
|
|
export class CollaborationRevisionGapError extends Error {
|
|
constructor(expected, actual) {
|
|
super(`Collaborative revision gap: expected ${expected}, received ${actual}`);
|
|
this.name = "CollaborationRevisionGapError";
|
|
this.expected = expected;
|
|
this.actual = actual;
|
|
}
|
|
}
|
|
|
|
export class CollaborationSession {
|
|
constructor(clientId) {
|
|
this.clientId = String(clientId || "");
|
|
this.ready = false;
|
|
this.serverContent = "";
|
|
this.serverOwnerMap = "[]";
|
|
this.revisionId = 0;
|
|
this.outstanding = null;
|
|
this.buffer = null;
|
|
this.nextUpdateId = 1;
|
|
}
|
|
|
|
initialize(content, ownerMap, revisionId, { clearPending = true } = {}) {
|
|
this.serverContent = String(content || "");
|
|
this.serverOwnerMap = ownerMap == null ? "[]" : String(ownerMap);
|
|
this.revisionId = Number(revisionId) || 0;
|
|
if (clearPending) {
|
|
this.outstanding = null;
|
|
this.buffer = null;
|
|
}
|
|
this.ready = true;
|
|
}
|
|
|
|
localDocument() {
|
|
return documentAfterPending(
|
|
this.serverContent,
|
|
this.serverOwnerMap,
|
|
this.outstanding,
|
|
this.buffer,
|
|
);
|
|
}
|
|
|
|
hasPending() {
|
|
return Boolean(this.outstanding || this.buffer);
|
|
}
|
|
|
|
queue(operation, ownerReplacements = []) {
|
|
operation = normalizeOperation(operation);
|
|
const replacements = mergeOwnerReplacements([], ownerReplacements);
|
|
if (!hasTextEffect(operation) && !replacements.length) return null;
|
|
|
|
const localLength = this.localDocument().content.length;
|
|
if (operationBaseLength(operation) !== localLength) {
|
|
throw new Error("Local operation base length does not match the collaborative document");
|
|
}
|
|
|
|
if (!this.buffer) {
|
|
this.buffer = pendingEnvelope(this.clientId, this.nextUpdateId++, operation, replacements);
|
|
} else {
|
|
this.buffer.operation = composeOperations(this.buffer.operation, operation);
|
|
this.buffer.ownerReplacements = mergeOwnerReplacements(this.buffer.ownerReplacements, replacements);
|
|
}
|
|
|
|
if (operationTargetLength(this.buffer.operation) !== this.localDocument().content.length) {
|
|
throw new Error("Buffered operation target length does not match the collaborative document");
|
|
}
|
|
return this.buffer;
|
|
}
|
|
|
|
sendable() {
|
|
if (!this.ready || this.outstanding || !this.buffer) return null;
|
|
return {
|
|
baseRevisionId: this.revisionId,
|
|
updateId: this.buffer.updateId,
|
|
operation: this.buffer.operation,
|
|
ownerReplacements: this.buffer.ownerReplacements,
|
|
};
|
|
}
|
|
|
|
markSent(updateId) {
|
|
if (this.outstanding || !this.buffer || this.buffer.updateId !== Number(updateId)) return false;
|
|
this.outstanding = this.buffer;
|
|
this.buffer = null;
|
|
return true;
|
|
}
|
|
|
|
integrate(message) {
|
|
const remote = serverEnvelope(message);
|
|
if (!Number.isFinite(remote.baseRevisionId) || !Number.isFinite(remote.revisionId)) {
|
|
throw new Error("Collaborative update is missing revision metadata");
|
|
}
|
|
if (remote.revisionId <= this.revisionId) return { duplicate: true, ownAck: false, remote: false };
|
|
if (remote.baseRevisionId !== this.revisionId) {
|
|
throw new CollaborationRevisionGapError(this.revisionId, remote.baseRevisionId);
|
|
}
|
|
|
|
const nextServer = applyOperationToDocument(
|
|
this.serverContent,
|
|
this.serverOwnerMap,
|
|
remote.operation,
|
|
remote.ownerReplacements,
|
|
);
|
|
const ownAck = Boolean(
|
|
this.outstanding
|
|
&& remote.clientId === this.clientId
|
|
&& remote.updateId === this.outstanding.updateId
|
|
);
|
|
|
|
if (ownAck) {
|
|
this.outstanding = null;
|
|
} else {
|
|
let remoteForPending = remote.operation;
|
|
if (this.outstanding) {
|
|
const outstandingBeforeRemote = compareOperationKeys(this.outstanding, remote) < 0;
|
|
const [outstandingPrime, remotePrime] = transformOperations(
|
|
this.outstanding.operation,
|
|
remoteForPending,
|
|
outstandingBeforeRemote,
|
|
);
|
|
this.outstanding.operation = outstandingPrime;
|
|
remoteForPending = remotePrime;
|
|
}
|
|
if (this.buffer) {
|
|
const bufferBeforeRemote = compareOperationKeys(this.buffer, remote) < 0;
|
|
const [bufferPrime] = transformOperations(
|
|
this.buffer.operation,
|
|
remoteForPending,
|
|
bufferBeforeRemote,
|
|
);
|
|
this.buffer.operation = bufferPrime;
|
|
}
|
|
}
|
|
|
|
this.serverContent = nextServer.content;
|
|
this.serverOwnerMap = nextServer.ownerMap;
|
|
this.revisionId = remote.revisionId;
|
|
return { duplicate: false, ownAck, remote: !ownAck };
|
|
}
|
|
|
|
resynchronize(message) {
|
|
const canonicalContent = String(message?.content || "");
|
|
const canonicalOwnerMap = message?.owner_map == null ? "[]" : String(message.owner_map);
|
|
const canonicalRevisionId = Number(message?.revision_id) || 0;
|
|
const outstanding = this.outstanding;
|
|
const buffer = this.buffer;
|
|
const local = this.localDocument();
|
|
const acknowledgedIds = (message?.acknowledged_update_ids || [])
|
|
.map(Number)
|
|
.filter(Number.isFinite);
|
|
const acknowledgedThrough = acknowledgedIds.length ? Math.max(...acknowledgedIds) : 0;
|
|
const messageClientId = String(message?.client_id ?? message?.clientId ?? "");
|
|
const messageUpdateId = Number(message?.update_id ?? message?.updateId ?? 0);
|
|
const outstandingAcknowledged = Boolean(
|
|
outstanding
|
|
&& (
|
|
acknowledgedThrough >= outstanding.updateId
|
|
|| (messageClientId === this.clientId && messageUpdateId === outstanding.updateId)
|
|
)
|
|
);
|
|
const pendingOwnerReplacements = mergeOwnerReplacements(
|
|
outstanding?.ownerReplacements || [],
|
|
buffer?.ownerReplacements || [],
|
|
);
|
|
let replayOperation = null;
|
|
let replayOwnerReplacements = [];
|
|
|
|
if (this.ready && (outstanding || buffer)) {
|
|
const canonicalAuthorship = parseAuthorship(canonicalContent, canonicalOwnerMap);
|
|
if (outstandingAcknowledged) {
|
|
if (buffer) {
|
|
const afterOutstanding = applyOperationToDocument(
|
|
this.serverContent,
|
|
this.serverOwnerMap,
|
|
outstanding.operation,
|
|
outstanding.ownerReplacements,
|
|
);
|
|
const missedRemote = operationFromEdit(
|
|
afterOutstanding.content,
|
|
canonicalContent,
|
|
canonicalAuthorship,
|
|
);
|
|
[replayOperation] = transformOperations(buffer.operation, missedRemote, false);
|
|
replayOwnerReplacements = buffer.ownerReplacements;
|
|
}
|
|
} else {
|
|
const localAuthorship = parseAuthorship(local.content, local.ownerMap);
|
|
const localOperation = operationFromEdit(this.serverContent, local.content, localAuthorship);
|
|
const missedRemote = operationFromEdit(
|
|
this.serverContent,
|
|
canonicalContent,
|
|
canonicalAuthorship,
|
|
);
|
|
[replayOperation] = transformOperations(localOperation, missedRemote, false);
|
|
replayOwnerReplacements = pendingOwnerReplacements;
|
|
}
|
|
}
|
|
|
|
this.initialize(canonicalContent, canonicalOwnerMap, canonicalRevisionId);
|
|
if (replayOperation || replayOwnerReplacements.length) {
|
|
this.queue(
|
|
replayOperation || identityOperation(canonicalContent.length),
|
|
replayOwnerReplacements,
|
|
);
|
|
}
|
|
return {
|
|
replayed: Boolean(this.buffer),
|
|
outstandingAcknowledged,
|
|
};
|
|
}
|
|
|
|
adoptCanonicalSnapshot(content, ownerMap, revisionId) {
|
|
revisionId = Number(revisionId) || 0;
|
|
if (revisionId !== this.revisionId) {
|
|
throw new CollaborationRevisionGapError(this.revisionId, revisionId);
|
|
}
|
|
if (String(content || "") !== this.serverContent) {
|
|
throw new Error("Canonical collaborative content does not match applied operations");
|
|
}
|
|
this.serverOwnerMap = ownerMap == null ? "[]" : String(ownerMap);
|
|
}
|
|
}
|