diff --git a/examples/excalidraw-example/src/hooks/useLoroSync.ts b/examples/excalidraw-example/src/hooks/useLoroSync.ts index ba48406..d881b04 100644 --- a/examples/excalidraw-example/src/hooks/useLoroSync.ts +++ b/examples/excalidraw-example/src/hooks/useLoroSync.ts @@ -1,9 +1,19 @@ import { useEffect, useRef, useCallback, useState } from "react"; import { throttle } from "throttle-debounce"; // TODO: REVIEW [stability] replace custom throttle with lib -import { LoroDoc, EphemeralStore, LoroEventBatch, LoroMap } from "loro-crdt"; +import { + LoroDoc, + EphemeralStore, + LoroEventBatch, + LoroMap, + type Value, +} from "loro-crdt"; import { LoroWebsocketClient } from "loro-websocket/client"; import { LoroAdaptor, LoroEphemeralAdaptor } from "loro-adaptors"; -import { ExcalidrawImperativeAPI } from "@excalidraw/excalidraw/types/types"; +import type { + AppState as ExcalidrawAppState, + ExcalidrawImperativeAPI, +} from "@excalidraw/excalidraw/types/types"; +import type { ExcalidrawElement } from "@excalidraw/excalidraw/types/element/types"; interface UseLoroSyncOptions { roomId: string; @@ -14,35 +24,28 @@ interface UseLoroSyncOptions { excalidrawAPI: React.RefObject; } -interface Collaborator { +interface PresenceEntry extends Record { userId: string; userName: string; userColor: string; - cursor?: { x: number; y: number }; + cursor?: CursorPosition; selectedElementIds?: string[]; lastActive: number; } -interface CursorPosition { +interface CursorPosition extends Record { x: number; y: number; } +interface Collaborator extends PresenceEntry {} +type AppState = ExcalidrawAppState; -// Minimal type definitions for Excalidraw (to avoid import issues) -export interface ExcalidrawElement { - id: string; - type: string; - x: number; - y: number; - width: number; - height: number; - version: number; - [key: string]: any; -} - -export interface AppState { - [key: string]: any; -} +type SceneUpdateArgs = Parameters< + ExcalidrawImperativeAPI["updateScene"] +>[0]; +type SceneElements = NonNullable; +type SceneAppStateUpdate = NonNullable; +type PresenceStoreState = Record; export function useLoroSync({ roomId, @@ -54,7 +57,8 @@ export function useLoroSync({ }: UseLoroSyncOptions) { const docRef = useRef(null); const clientRef = useRef(null); - const ephemeralRef = useRef> | null>(null); + const ephemeralRef = + useRef | null>(null); const [isConnected, setIsConnected] = useState(false); const [collaborators, setCollaborators] = useState>(new Map()); @@ -65,7 +69,7 @@ export function useLoroSync({ useEffect(() => { const doc = new LoroDoc(); const client = new LoroWebsocketClient({ url: wsUrl }); - const ephemeral = new EphemeralStore>(30000); // 30 second timeout + const ephemeral = new EphemeralStore(30000); // 30 second timeout docRef.current = doc; clientRef.current = client; @@ -81,7 +85,8 @@ export function useLoroSync({ if (event.by !== "local") { // Build scene data from doc and apply to Excalidraw. Avoid echo via flag. // TODO: REVIEW [avoid echo] We set a guard so the next Excalidraw onChange from updateScene is ignored. - const newElements = (elementsContainer.toJSON() || []) as ExcalidrawElement[]; + const newElements = + (elementsContainer.toJSON() || []) as SceneElements; const newAppState: Partial = {}; for (const [key, value] of appStateContainer.entries()) { newAppState[key as keyof AppState] = value; @@ -90,7 +95,11 @@ export function useLoroSync({ // Update checksum to match scene state const checksum = newElements.reduce((acc, e) => acc + (e?.version || 0), 0); lastChecksumRef.current = checksum; - excalidrawAPI.current?.updateScene({ elements: newElements as any, appState: newAppState as any }); + const sceneUpdate: SceneUpdateArgs = { elements: newElements }; + if (Object.keys(newAppState).length > 0) { + sceneUpdate.appState = newAppState as SceneAppStateUpdate; + } + excalidrawAPI.current?.updateScene(sceneUpdate); } }); @@ -197,6 +206,17 @@ export function useLoroSync({ const doc = docRef.current; const list = doc.getList("elements"); + const getMapAt = (index: number): LoroMap | undefined => { + const value = list.get(index); + return value instanceof LoroMap ? value : undefined; + }; + const ensureMapAt = (index: number): LoroMap => { + const map = getMapAt(index); + if (!map) { + throw new Error(`Expected LoroMap at index ${index}`); + } + return map; + }; // Filter out deleted const filtered = elements.filter(e => !e.isDeleted); @@ -205,9 +225,9 @@ export function useLoroSync({ const buildIndex = () => { const idx = new Map(); for (let i = 0; i < list.length; i++) { - const m = list.get(i) as unknown as LoroMap | undefined; - if (!m) continue; - const id = m.get("id") as string | undefined; + const map = getMapAt(i); + if (!map) continue; + const id = map.get("id") as string | undefined; if (id) idx.set(id, i); } return idx; @@ -223,9 +243,9 @@ export function useLoroSync({ if (pos == null) { // New element: insert at the desired position list.insertContainer(i, new LoroMap()); - const m = list.get(i) as unknown as LoroMap; + const map = ensureMapAt(i); for (const [k, v] of Object.entries(target)) { - m.set(k, v); + map.set(k, v); } changed = true; indexMap = buildIndex(); @@ -237,9 +257,9 @@ export function useLoroSync({ list.delete(pos, 1); const adjI = pos < i ? i - 1 : i; list.insertContainer(adjI, new LoroMap()); - const m = list.get(adjI) as unknown as LoroMap; + const map = ensureMapAt(adjI); for (const [k, v] of Object.entries(target)) { - m.set(k, v); + map.set(k, v); } changed = true; indexMap = buildIndex(); @@ -247,11 +267,11 @@ export function useLoroSync({ } // Same position: update only if version changed - const m = list.get(i) as unknown as LoroMap; - const prevVersion = m.get("version"); + const map = ensureMapAt(i); + const prevVersion = map.get("version"); if (prevVersion !== target.version) { for (const [k, v] of Object.entries(target)) { - m.set(k, v); + map.set(k, v); } changed = true; } diff --git a/examples/excalidraw-example/tsconfig.node.json b/examples/excalidraw-example/tsconfig.node.json index 099658c..97ede7e 100644 --- a/examples/excalidraw-example/tsconfig.node.json +++ b/examples/excalidraw-example/tsconfig.node.json @@ -4,7 +4,8 @@ "skipLibCheck": true, "module": "ESNext", "moduleResolution": "bundler", - "allowSyntheticDefaultImports": true + "allowSyntheticDefaultImports": true, + "strict": true }, "include": ["vite.config.ts"] -} \ No newline at end of file +} diff --git a/packages/loro-adaptors/src/elo-loro-adaptor.ts b/packages/loro-adaptors/src/elo-loro-adaptor.ts index a27bb76..f0542ae 100644 --- a/packages/loro-adaptors/src/elo-loro-adaptor.ts +++ b/packages/loro-adaptors/src/elo-loro-adaptor.ts @@ -110,8 +110,11 @@ export class EloLoroAdaptor implements CrdtDocAdaptor { if (spans.length === 1) { const { keyId, key } = await this.config.getPrivateKey(); - const s = spans[0]!; - const peerIdBytes = new TextEncoder().encode(String(s.peer)); + const [span] = spans; + if (!span) { + throw new Error("Expected delta span when packaging single update"); + } + const peerIdBytes = new TextEncoder().encode(String(span.peer)); const iv = this.config.ivFactory ? this.config.ivFactory() : undefined; @@ -119,8 +122,8 @@ export class EloLoroAdaptor implements CrdtDocAdaptor { updates, { peerId: peerIdBytes, - start: s.start, - end: s.start + s.length, + start: span.start, + end: span.start + span.length, keyId, iv, }, @@ -247,10 +250,11 @@ export class EloLoroAdaptor implements CrdtDocAdaptor { const mode = "snapshot"; const plaintext = this.doc.export({ mode }); const vvObj = vvToObject(this.doc.version()); + const encoder = new TextEncoder(); const vvEntries: Array<{ peerId: Uint8Array; counter: number }> = Object.keys(vvObj).map(peer => ({ - peerId: new TextEncoder().encode(peer), - counter: vvObj[peer]!, + peerId: encoder.encode(peer), + counter: vvObj[peer], })); const iv = this.config.ivFactory ? this.config.ivFactory() : undefined; const { record } = await encryptSnapshot( diff --git a/packages/loro-adaptors/src/index.ts b/packages/loro-adaptors/src/index.ts index 0a13669..c2240a0 100644 --- a/packages/loro-adaptors/src/index.ts +++ b/packages/loro-adaptors/src/index.ts @@ -1,2 +1,3 @@ export * from "./types"; export * from "./adaptors"; +export * from "./server"; diff --git a/packages/loro-adaptors/src/server/index.ts b/packages/loro-adaptors/src/server/index.ts new file mode 100644 index 0000000..4ce79fa --- /dev/null +++ b/packages/loro-adaptors/src/server/index.ts @@ -0,0 +1,4 @@ +export * from "./server-registry"; +export * from "./server-loro-adaptor"; +export * from "./server-loro-ephemeral-adaptor"; +export * from "./server-yjs-awareness-adaptor"; diff --git a/packages/loro-adaptors/src/server/server-loro-adaptor.ts b/packages/loro-adaptors/src/server/server-loro-adaptor.ts new file mode 100644 index 0000000..1b546e2 --- /dev/null +++ b/packages/loro-adaptors/src/server/server-loro-adaptor.ts @@ -0,0 +1,167 @@ +import { LoroDoc, VersionVector } from "loro-crdt"; +import { + CrdtType, + Permission, + MessageType, + JoinResponseOk, + UpdateError, + UpdateErrorCode, +} from "loro-protocol"; +import type { CrdtServerAdaptor } from "../types"; + +export class LoroServerAdaptor implements CrdtServerAdaptor { + readonly crdtType = CrdtType.Loro; + + createEmpty(): Uint8Array { + const doc = new LoroDoc(); + const snapshot = doc.export({ mode: "snapshot" }); + doc.free(); + return snapshot; + } + + handleJoinRequest( + documentData: Uint8Array, + clientVersion: Uint8Array, + permission: Permission + ): { + response: JoinResponseOk; + updates?: Uint8Array[]; + } { + const doc = new LoroDoc(); + try { + if (documentData.length > 0) { + doc.import(documentData); + } + + const serverVersion = doc.version(); + let updates: Uint8Array[] | undefined; + + if (clientVersion.length > 0) { + try { + const clientVV = VersionVector.decode(clientVersion); + const comparison = serverVersion.compare(clientVV); + + if (comparison && comparison > 0) { + const updateData = doc.export({ + mode: "update", + from: clientVV, + }); + updates = [updateData]; + } + } catch { + const snapshot = doc.export({ mode: "snapshot" }); + updates = [snapshot]; + } + } else { + const snapshot = doc.export({ mode: "snapshot" }); + updates = [snapshot]; + } + + const response: JoinResponseOk = { + type: MessageType.JoinResponseOk, + crdt: this.crdtType, + roomId: "", + permission, + version: serverVersion.encode(), + }; + + return { response, updates }; + } finally { + doc.free(); + } + } + + applyUpdates( + documentData: Uint8Array, + updates: Uint8Array[], + permission: Permission + ): { + success: boolean; + newDocumentData?: Uint8Array; + error?: UpdateError; + broadcastUpdates?: Uint8Array[]; + } { + if (permission === "read") { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.PermissionDenied, + message: "Read-only permission, cannot apply updates", + }, + }; + } + const doc = new LoroDoc(); + const broadcastUpdates: Uint8Array[] = []; + + try { + if (documentData.length > 0) { + doc.import(documentData); + } + for (const update of updates) { + if (update.length > 0) { + const importResult = doc.import(update); + if (importResult.success) { + broadcastUpdates.push(update); + } + } + } + + const newDocumentData = doc.export({ mode: "snapshot" }); + + return { + success: true, + newDocumentData, + broadcastUpdates: + broadcastUpdates.length > 0 ? broadcastUpdates : undefined, + }; + } catch (error) { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.InvalidUpdate, + message: error instanceof Error ? error.message : "Invalid update", + }, + }; + } finally { + doc.free(); + } + } + + getVersion(documentData: Uint8Array): Uint8Array { + const doc = new LoroDoc(); + try { + if (documentData.length > 0) { + doc.import(documentData); + } + return doc.version().encode(); + } finally { + doc.free(); + } + } + + getSize(documentData: Uint8Array): number { + return documentData.length; + } + + merge(documents: Uint8Array[]): Uint8Array { + const doc = new LoroDoc(); + try { + for (const docData of documents) { + if (docData.length > 0) { + doc.import(docData); + } + } + return doc.export({ mode: "snapshot" }); + } finally { + doc.free(); + } + } +} + +export const loroServerAdaptor = new LoroServerAdaptor(); diff --git a/packages/loro-adaptors/src/server/server-loro-ephemeral-adaptor.ts b/packages/loro-adaptors/src/server/server-loro-ephemeral-adaptor.ts new file mode 100644 index 0000000..21983db --- /dev/null +++ b/packages/loro-adaptors/src/server/server-loro-ephemeral-adaptor.ts @@ -0,0 +1,139 @@ +import { EphemeralStore } from "loro-crdt"; +import { + CrdtType, + Permission, + MessageType, + JoinResponseOk, + UpdateError, + UpdateErrorCode, +} from "loro-protocol"; +import type { CrdtServerAdaptor } from "../types"; + +export interface LoroEphemeralServerAdaptorConfig { + timeout?: number; +} + +export class LoroEphemeralServerAdaptor implements CrdtServerAdaptor { + readonly crdtType = CrdtType.LoroEphemeralStore; + private readonly timeout: number; + + constructor(config: LoroEphemeralServerAdaptorConfig = {}) { + this.timeout = config.timeout ?? 10_000; + } + + createEmpty(): Uint8Array { + const store = new EphemeralStore(this.timeout); + try { + return store.encodeAll(); + } finally { + store.inner.free(); + } + } + + handleJoinRequest( + documentData: Uint8Array, + _clientVersion: Uint8Array, + permission: Permission + ): { + response: JoinResponseOk; + updates?: Uint8Array[]; + } { + const response: JoinResponseOk = { + type: MessageType.JoinResponseOk, + crdt: this.crdtType, + roomId: "", + permission, + version: new Uint8Array(), + }; + + const updates = documentData.length > 0 ? [documentData] : undefined; + return { response, updates }; + } + + applyUpdates( + documentData: Uint8Array, + updates: Uint8Array[], + permission: Permission + ): { + success: boolean; + newDocumentData?: Uint8Array; + error?: UpdateError; + broadcastUpdates?: Uint8Array[]; + } { + if (permission === "read") { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.PermissionDenied, + message: "Read-only permission, cannot apply updates", + }, + }; + } + + const store = new EphemeralStore(this.timeout); + const broadcastUpdates: Uint8Array[] = []; + + try { + if (documentData.length > 0) { + store.apply(documentData); + } + for (const update of updates) { + if (update.length > 0) { + store.apply(update); + broadcastUpdates.push(update); + } + } + + const newDocumentData = store.encodeAll(); + + return { + success: true, + newDocumentData, + broadcastUpdates: + broadcastUpdates.length > 0 ? broadcastUpdates : undefined, + }; + } catch (error) { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.InvalidUpdate, + message: error instanceof Error ? error.message : "Invalid update", + }, + }; + } finally { + store.destroy(); + store.inner.free(); + } + } + + getVersion(_documentData: Uint8Array): Uint8Array { + return new Uint8Array(); + } + + getSize(documentData: Uint8Array): number { + return documentData.length; + } + + merge(documents: Uint8Array[]): Uint8Array { + const store = new EphemeralStore(this.timeout); + for (const data of documents) { + if (data.length > 0) { + store.apply(data); + } + } + try { + return store.encodeAll(); + } finally { + store.destroy(); + store.inner.free(); + } + } +} + +export const loroEphemeralServerAdaptor = new LoroEphemeralServerAdaptor(); diff --git a/packages/loro-adaptors/src/server/server-registry.ts b/packages/loro-adaptors/src/server/server-registry.ts new file mode 100644 index 0000000..d5d989b --- /dev/null +++ b/packages/loro-adaptors/src/server/server-registry.ts @@ -0,0 +1,50 @@ +import { CrdtType } from "loro-protocol"; +import type { AdaptorsForServer, CrdtServerAdaptor } from "../types"; + +class InMemoryAdaptorsForServer implements AdaptorsForServer { + private readonly adaptors = new Map(); + + register(adaptor: CrdtServerAdaptor): void { + this.adaptors.set(adaptor.crdtType, adaptor); + } + + registerMany(adaptors: Iterable): void { + for (const adaptor of adaptors) { + this.register(adaptor); + } + } + + get(crdtType: CrdtType): CrdtServerAdaptor | undefined { + return this.adaptors.get(crdtType); + } + + clear(): void { + this.adaptors.clear(); + } + + list(): CrdtServerAdaptor[] { + return Array.from(this.adaptors.values()); + } +} + +export const adaptorsForServer: AdaptorsForServer = new InMemoryAdaptorsForServer(); + +export function registerServerAdaptor(adaptor: CrdtServerAdaptor): void { + adaptorsForServer.register(adaptor); +} + +export function registerServerAdaptors(adaptors: Iterable): void { + adaptorsForServer.registerMany(adaptors); +} + +export function getServerAdaptor(crdtType: CrdtType): CrdtServerAdaptor | undefined { + return adaptorsForServer.get(crdtType); +} + +export function clearServerAdaptors(): void { + adaptorsForServer.clear(); +} + +export function listServerAdaptors(): CrdtServerAdaptor[] { + return adaptorsForServer.list(); +} diff --git a/packages/loro-adaptors/src/server/server-yjs-awareness-adaptor.ts b/packages/loro-adaptors/src/server/server-yjs-awareness-adaptor.ts new file mode 100644 index 0000000..78b3ff8 --- /dev/null +++ b/packages/loro-adaptors/src/server/server-yjs-awareness-adaptor.ts @@ -0,0 +1,120 @@ +import { + CrdtType, + Permission, + MessageType, + JoinResponseOk, + UpdateError, + UpdateErrorCode, +} from "loro-protocol"; +import type { CrdtServerAdaptor } from "../types"; + +export class YjsAwarenessServerAdaptor implements CrdtServerAdaptor { + readonly crdtType = CrdtType.YjsAwareness; + + createEmpty(): Uint8Array { + return new Uint8Array(); + } + + handleJoinRequest( + documentData: Uint8Array, + _clientVersion: Uint8Array, + permission: Permission + ): { + response: JoinResponseOk; + updates?: Uint8Array[]; + } { + const updates = documentData.length > 0 ? [documentData] : undefined; + + const response: JoinResponseOk = { + type: MessageType.JoinResponseOk, + crdt: this.crdtType, + roomId: "", + permission, + version: new Uint8Array(), + }; + + return { response, updates }; + } + + applyUpdates( + documentData: Uint8Array, + updates: Uint8Array[], + permission: Permission + ): { + success: boolean; + newDocumentData?: Uint8Array; + error?: UpdateError; + broadcastUpdates?: Uint8Array[]; + } { + if (permission === "read") { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.PermissionDenied, + message: "Read-only permission, cannot apply updates", + }, + }; + } + + try { + let total = documentData.length; + for (const u of updates) total += u.length; + const merged = new Uint8Array(total); + let offset = 0; + if (documentData.length) { + merged.set(documentData, 0); + offset += documentData.length; + } + for (const u of updates) { + if (u.length) { + merged.set(u, offset); + offset += u.length; + } + } + + return { + success: true, + newDocumentData: merged, + broadcastUpdates: updates.length > 0 ? updates : undefined, + }; + } catch (error) { + return { + success: false, + error: { + type: MessageType.UpdateError, + crdt: this.crdtType, + roomId: "", + code: UpdateErrorCode.InvalidUpdate, + message: error instanceof Error ? error.message : "Invalid update", + }, + }; + } + } + + getVersion(_documentData: Uint8Array): Uint8Array { + return new Uint8Array(); + } + + getSize(documentData: Uint8Array): number { + return documentData.length; + } + + merge(documents: Uint8Array[]): Uint8Array { + let total = 0; + for (const d of documents) total += d.length; + const out = new Uint8Array(total); + let offset = 0; + for (const d of documents) { + if (d.length) { + out.set(d, offset); + offset += d.length; + } + } + return out; + } +} + +export const yjsAwarenessServerAdaptor = new YjsAwarenessServerAdaptor(); diff --git a/packages/loro-adaptors/src/types.ts b/packages/loro-adaptors/src/types.ts index 86f63e6..3fc69d4 100644 --- a/packages/loro-adaptors/src/types.ts +++ b/packages/loro-adaptors/src/types.ts @@ -3,6 +3,7 @@ import { JoinResponseOk, JoinError, UpdateError, + Permission, } from "loro-protocol"; export interface CrdtDocAdaptor { @@ -58,3 +59,43 @@ export interface CrdtAdaptorContext { onJoinFailed: (reason: string) => void; onImportError: (error: Error, data: Uint8Array[]) => void; } + +export interface CrdtServerAdaptor { + readonly crdtType: CrdtType; + + createEmpty(): Uint8Array; + + handleJoinRequest( + documentData: Uint8Array, + clientVersion: Uint8Array, + permission: Permission + ): { + response: JoinResponseOk; + updates?: Uint8Array[]; + }; + + applyUpdates( + documentData: Uint8Array, + updates: Uint8Array[], + permission: Permission + ): { + success: boolean; + newDocumentData?: Uint8Array; + error?: UpdateError; + broadcastUpdates?: Uint8Array[]; + }; + + getVersion(documentData: Uint8Array): Uint8Array; + + getSize(documentData: Uint8Array): number; + + merge(documents: Uint8Array[]): Uint8Array; +} + +export interface AdaptorsForServer { + register(adaptor: CrdtServerAdaptor): void; + registerMany(adaptors: Iterable): void; + get(crdtType: CrdtType): CrdtServerAdaptor | undefined; + clear(): void; + list(): CrdtServerAdaptor[]; +} diff --git a/packages/loro-adaptors/tests/elo-adaptor.test.ts b/packages/loro-adaptors/tests/elo-adaptor.test.ts index 9770e11..cb08c41 100644 --- a/packages/loro-adaptors/tests/elo-adaptor.test.ts +++ b/packages/loro-adaptors/tests/elo-adaptor.test.ts @@ -51,7 +51,9 @@ describe("EloLoroAdaptor — snapshot join", () => { const [container] = updates; const records = decodeEloContainer(container); expect(records.length).toBe(1); - const parsed = parseEloRecordHeader(records[0]!); + const record = records[0]; + if (!record) throw new Error("Expected an ELO record to be emitted"); + const parsed = parseEloRecordHeader(record); expect(parsed.kind).toBe(0x01); // Snapshot expect(parsed.keyId).toBe("k1"); expect(parsed.iv.length).toBe(12); diff --git a/packages/loro-protocol/tests/e2ee.test.ts b/packages/loro-protocol/tests/e2ee.test.ts index 401d9bf..c0108b0 100644 --- a/packages/loro-protocol/tests/e2ee.test.ts +++ b/packages/loro-protocol/tests/e2ee.test.ts @@ -15,8 +15,11 @@ describe("%ELO container codec", () => { const bytes = encodeEloContainer(records); const out = decodeEloContainer(bytes); expect(out.length).toBe(2); - expect(Array.from(out[0]!)).toEqual([1, 2, 3]); - expect(Array.from(out[1]!)).toEqual([4]); + const first = out[0]; + const second = out[1]; + if (!first || !second) throw new Error("Expected two decoded records"); + expect(Array.from(first)).toEqual([1, 2, 3]); + expect(Array.from(second)).toEqual([4]); }); }); diff --git a/packages/loro-protocol/tests/encoding.test.ts b/packages/loro-protocol/tests/encoding.test.ts index 5244c26..8bd2553 100644 --- a/packages/loro-protocol/tests/encoding.test.ts +++ b/packages/loro-protocol/tests/encoding.test.ts @@ -185,7 +185,7 @@ describe("Message Encoding and Decoding", () => { expect(decoded.type).toBe(MessageType.DocUpdate); if (decoded.type !== MessageType.DocUpdate) throw new Error("bad type"); expect(decoded.updates.length).toBe(1); - expect(Array.from(decoded.updates[0]!)).toEqual([111, 222]); + expect(Array.from(decoded.updates[0])).toEqual([111, 222]); }); it("encodes and decodes DocUpdate with multiple updates", () => { @@ -206,9 +206,9 @@ describe("Message Encoding and Decoding", () => { expect(decoded.type).toBe(MessageType.DocUpdate); if (decoded.type !== MessageType.DocUpdate) throw new Error("bad type"); expect(decoded.updates.length).toBe(3); - expect(Array.from(decoded.updates[0]!)).toEqual([1, 2, 3]); - expect(Array.from(decoded.updates[1]!)).toEqual([4, 5, 6, 7]); - expect(Array.from(decoded.updates[2]!)).toEqual([8]); + expect(Array.from(decoded.updates[0])).toEqual([1, 2, 3]); + expect(Array.from(decoded.updates[1])).toEqual([4, 5, 6, 7]); + expect(Array.from(decoded.updates[2])).toEqual([8]); }); it("encodes and decodes DocUpdate with empty updates array", () => { diff --git a/packages/loro-websocket/test-wrappers/recv-elo-doc.ts b/packages/loro-websocket/test-wrappers/recv-elo-doc.ts index c7123ed..26d6017 100644 --- a/packages/loro-websocket/test-wrappers/recv-elo-doc.ts +++ b/packages/loro-websocket/test-wrappers/recv-elo-doc.ts @@ -52,7 +52,10 @@ async function main() { process.exit(1); } -const isEntrypoint = import.meta.url === pathToFileURL(process.argv[1]!).href; +const scriptArg = process.argv[1]; +const isEntrypoint = + typeof scriptArg === "string" && + import.meta.url === pathToFileURL(scriptArg).href; if (isEntrypoint) { // eslint-disable-next-line unicorn/prefer-top-level-await main().catch(err => { diff --git a/packages/loro-websocket/test-wrappers/recv-index.ts b/packages/loro-websocket/test-wrappers/recv-index.ts index 9b43671..5a599a1 100644 --- a/packages/loro-websocket/test-wrappers/recv-index.ts +++ b/packages/loro-websocket/test-wrappers/recv-index.ts @@ -61,7 +61,10 @@ async function main() { } // ESM entrypoint guard -const isEntrypoint = import.meta.url === pathToFileURL(process.argv[1]!).href; +const scriptArg = process.argv[1]; +const isEntrypoint = + typeof scriptArg === "string" && + import.meta.url === pathToFileURL(scriptArg).href; if (isEntrypoint) { // eslint-disable-next-line unicorn/prefer-top-level-await main().catch(err => { diff --git a/packages/loro-websocket/test-wrappers/send-elo-normative.ts b/packages/loro-websocket/test-wrappers/send-elo-normative.ts index de3c198..53da4d6 100644 --- a/packages/loro-websocket/test-wrappers/send-elo-normative.ts +++ b/packages/loro-websocket/test-wrappers/send-elo-normative.ts @@ -56,7 +56,10 @@ async function main() { } // ESM entrypoint guard -const isEntrypoint = import.meta.url === pathToFileURL(process.argv[1]!).href; +const scriptArg = process.argv[1]; +const isEntrypoint = + typeof scriptArg === "string" && + import.meta.url === pathToFileURL(scriptArg).href; if (isEntrypoint) { // eslint-disable-next-line unicorn/prefer-top-level-await main().catch(err => { diff --git a/packages/loro-websocket/test-wrappers/start-simple-server.ts b/packages/loro-websocket/test-wrappers/start-simple-server.ts index 9cbe2c2..ac78b90 100644 --- a/packages/loro-websocket/test-wrappers/start-simple-server.ts +++ b/packages/loro-websocket/test-wrappers/start-simple-server.ts @@ -30,7 +30,10 @@ async function main() { } // ESM entrypoint guard -const isEntrypoint = import.meta.url === pathToFileURL(process.argv[1]!).href; +const scriptArg = process.argv[1]; +const isEntrypoint = + typeof scriptArg === "string" && + import.meta.url === pathToFileURL(scriptArg).href; if (isEntrypoint) { // eslint-disable-next-line unicorn/prefer-top-level-await main().catch(err => { diff --git a/packages/loro-websocket/tests/e2e-elo.test.ts b/packages/loro-websocket/tests/e2e-elo.test.ts index d98b114..4e6ef66 100644 --- a/packages/loro-websocket/tests/e2e-elo.test.ts +++ b/packages/loro-websocket/tests/e2e-elo.test.ts @@ -368,9 +368,9 @@ async function waitForUpdateError(ws: WebSocket): Promise { ws.on("message", (data: unknown) => { try { const msg = tryDecodeMsg(data); - if (msg && msg.type === MessageType.UpdateError) { + if (msg?.type === MessageType.UpdateError) { clearTimeout(t); - resolve(msg as UpdateError); + resolve(msg); } } catch {} }); diff --git a/packages/loro-websocket/tests/e2e.test.ts b/packages/loro-websocket/tests/e2e.test.ts index 6aafc82..3009e99 100644 --- a/packages/loro-websocket/tests/e2e.test.ts +++ b/packages/loro-websocket/tests/e2e.test.ts @@ -2,10 +2,74 @@ import { describe, it, expect, beforeAll, afterAll } from "vitest"; import { WebSocket } from "ws"; import getPort from "get-port"; import { SimpleServer } from "../src/server/simple-server"; -import { LoroWebsocketClient } from "../src/client"; -import { ClientStatus } from "../src/client"; +import { + registerCrdtDoc, + getCrdtDocConstructor, + type CrdtDoc, + type CrdtDocConstructor, +} from "../src/server/crdt-doc"; +import { LoroWebsocketClient, ClientStatus } from "../src/client"; import type { LoroWebsocketClientRoom } from "../src/client"; -import { createLoroAdaptor } from "loro-adaptors"; +import { createLoroAdaptor, loroServerAdaptor } from "loro-adaptors"; +import { CrdtType } from "loro-protocol"; + +class AdaptorBackedLoroDoc implements CrdtDoc { + private data: Uint8Array; + + constructor() { + this.data = AdaptorBackedLoroDoc.clone(loroServerAdaptor.createEmpty()); + } + + private static clone(input: Uint8Array): Uint8Array { + return input.length ? new Uint8Array(input) : new Uint8Array(); + } + + getVersion(): Uint8Array { + return AdaptorBackedLoroDoc.clone(loroServerAdaptor.getVersion(this.data)); + } + + computeBackfill(clientVersion: Uint8Array): Uint8Array[] | null { + const version = clientVersion ?? new Uint8Array(); + const { updates } = loroServerAdaptor.handleJoinRequest( + this.data, + version, + "write" + ); + return updates && updates.length + ? updates.map(update => AdaptorBackedLoroDoc.clone(update)) + : null; + } + + applyUpdates(updates: Uint8Array[]) { + const result = loroServerAdaptor.applyUpdates(this.data, updates, "write"); + if (!result.success) { + const message = result.error?.message ?? "Unknown update failure"; + return { ok: false as const, error: message }; + } + + if (result.newDocumentData) { + this.data = AdaptorBackedLoroDoc.clone(result.newDocumentData); + } + + return { ok: true as const }; + } + + shouldPersist(): boolean { + return true; + } + + exportSnapshot(): Uint8Array | null { + return this.data.length ? AdaptorBackedLoroDoc.clone(this.data) : null; + } + + importSnapshot(data: Uint8Array): void { + this.data = AdaptorBackedLoroDoc.clone(data); + } + + allowBackfillWhenNoOtherClients(): boolean { + return false; + } +} // Make WebSocket available globally for the client Object.defineProperty(globalThis, "WebSocket", { @@ -17,8 +81,12 @@ Object.defineProperty(globalThis, "WebSocket", { describe("E2E: Client-Server Sync", () => { let server: SimpleServer; let port: number; + let restoreLoroCtor: CrdtDocConstructor | undefined; beforeAll(async () => { + restoreLoroCtor = getCrdtDocConstructor(CrdtType.Loro); + registerCrdtDoc(CrdtType.Loro, () => new AdaptorBackedLoroDoc()); + port = await getPort(); server = new SimpleServer({ port }); await server.start(); @@ -26,6 +94,9 @@ describe("E2E: Client-Server Sync", () => { afterAll(async () => { await server.stop(); + if (restoreLoroCtor) { + registerCrdtDoc(CrdtType.Loro, restoreLoroCtor); + } }, 15000); it("should sync two clients through server", async () => { @@ -238,7 +309,9 @@ describe("E2E: Client-Server Sync", () => { // Getter should reflect last RTT const got = client.getLatency(); - expect(typeof got === "number" && isFinite(got!)).toBe(true); + const hasFiniteLatency = + typeof got === "number" && Number.isFinite(got); + expect(hasFiniteLatency).toBe(true); // Subscribe again; should emit immediately with current latency let immediate: number | undefined; @@ -580,11 +653,8 @@ function installMockWindow(initialOnline = true) { for (const listener of Array.from(set)) { if (typeof listener === "function") { listener.call(mockWindow, evt); - } else if ( - listener && - typeof (listener as EventListenerObject).handleEvent === "function" - ) { - (listener as EventListenerObject).handleEvent.call(mockWindow, evt); + } else if (isEventListenerObject(listener)) { + listener.handleEvent.call(mockWindow, evt); } } }; @@ -600,7 +670,7 @@ function installMockWindow(initialOnline = true) { if (originalWindowDescriptor) { Object.defineProperty(globalThis, "window", originalWindowDescriptor); } else { - delete (globalThis as any).window; + Reflect.deleteProperty(globalThis, "window"); } if (originalNavigatorDescriptor) { Object.defineProperty( @@ -609,13 +679,24 @@ function installMockWindow(initialOnline = true) { originalNavigatorDescriptor ); } else { - delete (globalThis as any).navigator; + Reflect.deleteProperty(globalThis, "navigator"); } listeners.clear(); }, }; } +function isEventListenerObject( + value: EventListenerOrEventListenerObject +): value is EventListenerObject { + return ( + typeof value === "object" && + value !== null && + "handleEvent" in value && + typeof value.handleEvent === "function" + ); +} + // Small polling helper for this file async function waitUntil( cond: () => boolean,