From d09b1dfc2e0821414d5bff7c687d1628549400e2 Mon Sep 17 00:00:00 2001 From: A0 Orchestrator Date: Tue, 30 Jun 2026 12:42:49 +0200 Subject: [PATCH] =?UTF-8?q?fix:=20HIGH=20issues=20#13,#14,#18,#19,#20=20?= =?UTF-8?q?=E2=80=94=20frontend=20state/sync:=20copy-paste=20selection,=20?= =?UTF-8?q?dynamic=20activeLayer,=20group=20creation,=20Yjs=20awareness=20?= =?UTF-8?q?protocol,=20connect=20sync?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/package-lock.json | 20 ++++ backend/package.json | 1 + backend/src/websocket/yjsServer.ts | 95 +++++++++++++-- frontend/src/App.tsx | 32 +++-- frontend/src/crdt/AwarenessManager.ts | 158 +++++++++++++++++++------ frontend/src/crdt/WebSocketProvider.ts | 96 +++++++++++++-- frontend/src/crdt/index.ts | 2 +- frontend/src/crdt/useYjsBinding.ts | 11 +- 8 files changed, 353 insertions(+), 62 deletions(-) diff --git a/backend/package-lock.json b/backend/package-lock.json index 821695a..3115236 100644 --- a/backend/package-lock.json +++ b/backend/package-lock.json @@ -15,6 +15,7 @@ "better-sqlite3": "^11.0.0", "fastify": "^4.28.0", "y-leveldb": "^0.2.0", + "y-protocols": "^1.0.7", "yjs": "^13.6.0" }, "devDependencies": { @@ -3868,6 +3869,25 @@ "yjs": "^13.0.0" } }, + "node_modules/y-protocols": { + "version": "1.0.7", + "resolved": "https://registry.npmjs.org/y-protocols/-/y-protocols-1.0.7.tgz", + "integrity": "sha512-YSVsLoXxO67J6eE/nV4AtFtT3QEotZf5sK5BHxFBXso7VDUT3Tx07IfA6hsu5Q5OmBdMkQVmFZ9QOA7fikWvnw==", + "dependencies": { + "lib0": "^0.2.85" + }, + "engines": { + "node": ">=16.0.0", + "npm": ">=8.0.0" + }, + "funding": { + "type": "GitHub Sponsors ❤", + "url": "https://github.com/sponsors/dmonad" + }, + "peerDependencies": { + "yjs": "^13.0.0" + } + }, "node_modules/yjs": { "version": "13.6.31", "resolved": "https://registry.npmjs.org/yjs/-/yjs-13.6.31.tgz", diff --git a/backend/package.json b/backend/package.json index 8e28745..8303263 100644 --- a/backend/package.json +++ b/backend/package.json @@ -18,6 +18,7 @@ "better-sqlite3": "^11.0.0", "fastify": "^4.28.0", "y-leveldb": "^0.2.0", + "y-protocols": "^1.0.7", "yjs": "^13.6.0" }, "devDependencies": { diff --git a/backend/src/websocket/yjsServer.ts b/backend/src/websocket/yjsServer.ts index 3bce38f..fc9935e 100644 --- a/backend/src/websocket/yjsServer.ts +++ b/backend/src/websocket/yjsServer.ts @@ -1,16 +1,27 @@ /** - * Yjs WebSocket Server – Real-time collaboration with LevelDB persistence + * Yjs WebSocket Server – Real-time collaboration with LevelDB persistence. + * + * Message protocol (byte-prefixed): + * 0 = Y.Doc update (Y.encodeStateAsUpdate or Y.encodeStateAsUpdateYUpdate) + * 1 = Awareness update (encodeAwarenessUpdate from y-protocols/awareness) */ import * as Y from 'yjs'; import { LeveldbPersistence } from 'y-leveldb'; +import { Awareness, encodeAwarenessUpdate, applyAwarenessUpdate, removeAwarenessStates } from 'y-protocols/awareness'; import type { FastifyInstance } from 'fastify'; import type { WebSocket } from '@fastify/websocket'; import type { AuthService } from '../auth/AuthService.js'; const PERSISTENCE_DIR = process.env.YJS_PERSISTENCE_DIR || '/tmp/yjs-documents'; +// Message type prefixes (must match frontend WebSocketProvider.ts) +const MSG_DOC_UPDATE = 0; +const MSG_AWARENESS = 1; + // Document store: docName → Y.Doc const docs = new Map(); +// Awareness store: docName → Awareness +const awarenessMap = new Map(); // Connection tracking: docName → Set const connections = new Map>(); let persistence: LeveldbPersistence | null = null; @@ -46,6 +57,10 @@ async function getOrCreateDoc(docName: string): Promise { const doc = new Y.Doc(); docs.set(docName, doc); + // Create awareness instance for this document + const awareness = new Awareness(doc); + awarenessMap.set(docName, awareness); + // Load persisted state if available if (persistence && persistenceReady) { try { @@ -72,14 +87,23 @@ async function getOrCreateDoc(docName: string): Promise { return doc; } -function broadcastUpdate(docName: string, update: Uint8Array, exclude?: WebSocket): void { +function getOrCreateAwareness(docName: string): Awareness | null { + return awarenessMap.get(docName) ?? null; +} + +function broadcastUpdate(docName: string, update: Uint8Array, msgType: number, exclude?: WebSocket): void { const conns = connections.get(docName); if (!conns) return; + // Prepend message type byte + const msg = new Uint8Array(update.length + 1); + msg[0] = msgType; + msg.set(update, 1); + for (const ws of conns) { if (ws !== exclude && ws.readyState === 1 /* OPEN */) { try { - ws.send(update); + ws.send(msg); } catch { // Connection error — will be cleaned up on close } @@ -108,6 +132,7 @@ export function registerYjsWebSocket(fastify: FastifyInstance, authService: Auth } const doc = await getOrCreateDoc(docName); + const awareness = getOrCreateAwareness(docName); // Track connection if (!connections.has(docName)) { @@ -115,18 +140,50 @@ export function registerYjsWebSocket(fastify: FastifyInstance, authService: Auth } connections.get(docName)!.add(socket); - // Send current document state to new client + // Send current document state to new client (with MSG_DOC_UPDATE prefix) const stateUpdate = Y.encodeStateAsUpdate(doc); if (stateUpdate.length > 0 && socket.readyState === 1) { - socket.send(stateUpdate); + const msg = new Uint8Array(stateUpdate.length + 1); + msg[0] = MSG_DOC_UPDATE; + msg.set(stateUpdate, 1); + socket.send(msg); } - // Listen for updates from this client + // Send current awareness states to new client + if (awareness && awareness.getStates().size > 0 && socket.readyState === 1) { + const awarenessUpdate = encodeAwarenessUpdate(awareness, Array.from(awareness.getStates().keys())); + if (awarenessUpdate.length > 0) { + const msg = new Uint8Array(awarenessUpdate.length + 1); + msg[0] = MSG_AWARENESS; + msg.set(awarenessUpdate, 1); + socket.send(msg); + } + } + + // Listen for messages from this client socket.on('message', (data: Buffer) => { try { - const update = new Uint8Array(data); - Y.applyUpdate(doc, update); - broadcastUpdate(docName, update, socket); + const raw = new Uint8Array(data); + if (raw.length === 0) return; + + const msgType = raw[0]; + const payload = raw.slice(1); + + if (msgType === MSG_AWARENESS) { + // Awareness update — apply to server awareness and broadcast + if (awareness) { + applyAwarenessUpdate(awareness, payload, socket); + broadcastUpdate(docName, payload, MSG_AWARENESS, socket); + } + } else if (msgType === MSG_DOC_UPDATE) { + // Y.Doc update — apply to server doc and broadcast + Y.applyUpdate(doc, payload); + broadcastUpdate(docName, payload, MSG_DOC_UPDATE, socket); + } else { + // Legacy: treat entire message as a Y.Doc update (no prefix) + Y.applyUpdate(doc, raw); + broadcastUpdate(docName, raw, MSG_DOC_UPDATE, socket); + } } catch (err) { console.warn(`Invalid update from client for ${docName}:`, err); } @@ -137,6 +194,26 @@ export function registerYjsWebSocket(fastify: FastifyInstance, authService: Auth const conns = connections.get(docName); if (conns) { conns.delete(socket); + + // Remove this client's awareness state + if (awareness) { + // Find and remove the client's awareness state based on socket origin + const statesToRemove: number[] = []; + for (const [clientID, meta] of awareness.meta) { + if (meta === socket) { + statesToRemove.push(clientID); + } + } + if (statesToRemove.length > 0) { + removeAwarenessStates(awareness, statesToRemove, socket); + // Broadcast removal to remaining clients + const removalUpdate = encodeAwarenessUpdate(awareness, statesToRemove); + if (removalUpdate.length > 0) { + broadcastUpdate(docName, removalUpdate, MSG_AWARENESS, socket); + } + } + } + if (conns.size === 0) { connections.delete(docName); // Keep doc in memory for reconnection (persistence handles durability) diff --git a/frontend/src/App.tsx b/frontend/src/App.tsx index 231e08a..9838d06 100644 --- a/frontend/src/App.tsx +++ b/frontend/src/App.tsx @@ -111,6 +111,7 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack // Right sidebar const [activeRightPanel, setActiveRightPanel] = useState('tool'); const [selectedElement, setSelectedElement] = useState(null); + const [selectedElementIds, setSelectedElementIds] = useState([]); const [layers, setLayers] = useState(initialLayers); const [activeLayerId, setActiveLayerId] = useState('layer-0'); const [blocks, setBlocks] = useState(initialBlocks); @@ -191,6 +192,8 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack bgConfigRef.current = bgConfig; const activeLayerIdRef = useRef(activeLayerId); activeLayerIdRef.current = activeLayerId; + const selectedElementIdsRef = useRef(selectedElementIds); + selectedElementIdsRef.current = selectedElementIds; // Clipboard (for copy/paste) const clipboardRef = React.useRef(null); @@ -642,11 +645,16 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack }, [blocks, activeLayerId, collab, drawingId, token]); const handleSelectionChange = useCallback((selectedIds: string[]) => { + setSelectedElementIds(selectedIds); if (selectedIds.length === 1) { const el = elements.find(e => e.id === selectedIds[0]); setSelectedElement(el ?? null); - } else { + } else if (selectedIds.length === 0) { setSelectedElement(null); + } else { + // Multiple elements selected — keep the first one as a representative for the properties panel + const el = elements.find(e => e.id === selectedIds[0]); + setSelectedElement(el ?? null); } }, [elements]); @@ -776,7 +784,8 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack if (action === 'undo') { handleUndo(); return; } if (action === 'redo') { handleRedo(); return; } if (action === 'copy') { - const selected = selectedElement ? [selectedElement] : []; + const ids = selectedElementIdsRef.current; + const selected = ids.length > 0 ? elements.filter(e => ids.includes(e.id)) : []; const clip = selected.length > 0 ? selected : elements.slice(0, 1); clipboardRef.current = clip; setCommandHistory((prev) => [...prev, { prefix: '·', text: `${clip.length} Element(e) kopiert`, type: 'info' }]); @@ -1136,10 +1145,17 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack // Group command: create group from currently selected elements if (upper === 'GROUP' || upper === 'GRP') { const gm = groupManagerRef.current; - // For now, group all elements (selection state is in InteractionEngine, not accessible here) - // In a full implementation, we'd need selected IDs from the interaction engine - setCommandHistory((prev) => [...prev, { prefix: '·', text: 'Gruppe erstellt (Auswahl im Canvas erforderlich)', type: 'info' }]); - setGroups(gm.getGroups()); + const ids = selectedElementIdsRef.current; + if (ids.length < 2) { + setCommandHistory((prev) => [...prev, { prefix: '·', text: 'Gruppe: Mindestens 2 Elemente auswählen', type: 'info' }]); + return; + } + const group = gm.createGroup(ids); + const newGroups = gm.getGroups(); + setGroups(newGroups); + // Sync group to Yjs CRDT for real-time collaboration + collab.setGroup(group); + setCommandHistory((prev) => [...prev, { prefix: '·', text: `Gruppe erstellt: ${group.name} (${ids.length} Elemente)`, type: 'info' }]); return; } @@ -1194,7 +1210,7 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack setCommandHistory((prev) => [...prev, { prefix: '·', text: `Unbekannter Befehl: ${cmd}`, type: 'info' }]); } } - }, [handleUndo, handleRedo, handleRibbonAction]); + }, [handleUndo, handleRedo, handleRibbonAction, collab]); const handleKISend = useCallback(async (text: string) => { const userMsg: KIMessage = { id: `ki-${Date.now()}`, role: 'user', content: text }; @@ -1249,7 +1265,7 @@ const CADEditor: React.FC = ({ projectId, token, onNavigateBack handleKISend(suggestion.label); }, [handleKISend]); - const activeLayerName = layers.find((l) => l.id === 'layer-0')?.name ?? '—'; + const activeLayerName = layers.find((l) => l.id === activeLayerId)?.name ?? '—'; // ─── Render ───────────────────────────────────────────── return ( diff --git a/frontend/src/crdt/AwarenessManager.ts b/frontend/src/crdt/AwarenessManager.ts index 021dcd2..02193c9 100644 --- a/frontend/src/crdt/AwarenessManager.ts +++ b/frontend/src/crdt/AwarenessManager.ts @@ -1,9 +1,11 @@ /** * AwarenessManager – Tracks user presence, cursor position, and selection. - * Uses a Y.Map inside the shared doc so awareness state syncs via the same - * raw-update WebSocket protocol as the rest of the document. + * Uses the Yjs awareness protocol (y-protocols/awareness) for ephemeral + * state that should NOT be persisted in the shared Y.Doc. + * + * Reference: https://docs.yjs.dev/ecosystem/awareness */ -import * as Y from 'yjs'; +import { Awareness, encodeAwarenessUpdate, applyAwarenessUpdate } from 'y-protocols/awareness'; import type { YjsDocument } from './YjsDocument'; export interface UserCursor { @@ -20,13 +22,28 @@ export interface UserSelection { elementIds: string[]; } -export interface AwarenessState { - cursors: Y.Map; - selections: Y.Map; +/** + * Payload stored under the 'cursor' field of each client's awareness state. + */ +interface CursorAwarenessPayload { + userId: string; + userName: string; + color: string; + x: number; + y: number; + visible: boolean; +} + +/** + * Payload stored under the 'selection' field of each client's awareness state. + */ +interface SelectionAwarenessPayload { + userId: string; + elementIds: string[]; } export class AwarenessManager { - readonly state: AwarenessState; + readonly awareness: Awareness; private yjsDoc: YjsDocument; private userId: string; private listeners: Set<() => void> = new Set(); @@ -34,85 +51,160 @@ export class AwarenessManager { constructor(yjsDoc: YjsDocument, userId: string) { this.yjsDoc = yjsDoc; this.userId = userId; - this.state = { - cursors: this.yjsDoc.doc.getMap('awareness-cursors'), - selections: this.yjsDoc.doc.getMap('awareness-selections'), - }; + this.awareness = new Awareness(yjsDoc.doc); } /** Set the local user's cursor position */ setCursor(x: number, y: number, userName: string, color: string): void { - this.state.cursors.set(this.userId, { + const payload: CursorAwarenessPayload = { userId: this.userId, userName, color, x, y, visible: true, - }); + }; + this.awareness.setLocalStateField('cursor', payload); } /** Hide the local user's cursor */ hideCursor(): void { - const cur = this.state.cursors.get(this.userId); - if (cur) { - this.state.cursors.set(this.userId, { ...cur, visible: false }); + const localState = this.awareness.getLocalState(); + if (localState && localState.cursor) { + const cursor = localState.cursor as CursorAwarenessPayload; + this.awareness.setLocalStateField('cursor', { + ...cursor, + visible: false, + } as CursorAwarenessPayload); } } /** Set the local user's selection */ setSelection(elementIds: string[]): void { - this.state.selections.set(this.userId, { + const payload: SelectionAwarenessPayload = { userId: this.userId, elementIds, - }); + }; + this.awareness.setLocalStateField('selection', payload); } /** Clear the local user's selection */ clearSelection(): void { - this.state.selections.delete(this.userId); + this.awareness.setLocalStateField('selection', null); } /** Remove the local user from awareness (on disconnect) */ removeSelf(): void { - this.state.cursors.delete(this.userId); - this.state.selections.delete(this.userId); + this.awareness.setLocalState(null); } - /** Get all active cursors */ + /** + * Encode an awareness update for the local client. + * Used by WebSocketProvider to send awareness state over the wire. + */ + encodeLocalState(): Uint8Array | null { + const localState = this.awareness.getLocalState(); + if (!localState) return null; + return encodeAwarenessUpdate(this.awareness, [this.awareness.clientID]); + } + + /** + * Apply a remote awareness update received from the server. + * Used by WebSocketProvider to process incoming awareness messages. + */ + applyRemoteUpdate(update: Uint8Array, origin: unknown = 'remote'): void { + applyAwarenessUpdate(this.awareness, update, origin); + this.notifyListeners(); + } + + /** Get all active cursors from all connected clients */ getAllCursors(): UserCursor[] { - return Array.from(this.state.cursors.values()).filter((c) => c.visible); + const cursors: UserCursor[] = []; + for (const [, state] of this.awareness.getStates()) { + if (state && state.cursor) { + const c = state.cursor as CursorAwarenessPayload; + if (c.visible) { + cursors.push({ + userId: c.userId, + userName: c.userName, + color: c.color, + x: c.x, + y: c.y, + visible: c.visible, + }); + } + } + } + return cursors; } - /** Get all selections */ + /** Get all selections from all connected clients */ getAllSelections(): UserSelection[] { - return Array.from(this.state.selections.values()); + const selections: UserSelection[] = []; + for (const [, state] of this.awareness.getStates()) { + if (state && state.selection) { + const s = state.selection as SelectionAwarenessPayload; + selections.push({ + userId: s.userId, + elementIds: s.elementIds, + }); + } + } + return selections; } /** Get a specific user's cursor */ getCursor(userId: string): UserCursor | undefined { - return this.state.cursors.get(userId); + for (const [, state] of this.awareness.getStates()) { + if (state && state.cursor) { + const c = state.cursor as CursorAwarenessPayload; + if (c.userId === userId) { + return { + userId: c.userId, + userName: c.userName, + color: c.color, + x: c.x, + y: c.y, + visible: c.visible, + }; + } + } + } + return undefined; } /** Get a specific user's selection */ getSelection(userId: string): UserSelection | undefined { - return this.state.selections.get(userId); + for (const [, state] of this.awareness.getStates()) { + if (state && state.selection) { + const s = state.selection as SelectionAwarenessPayload; + if (s.userId === userId) { + return { + userId: s.userId, + elementIds: s.elementIds, + }; + } + } + } + return undefined; } /** Register a callback for awareness changes */ onChange(callback: () => void): () => void { this.listeners.add(callback); - const cursorObserver = () => this.notifyListeners(); - const selectionObserver = () => this.notifyListeners(); - this.state.cursors.observe(cursorObserver); - this.state.selections.observe(selectionObserver); + const observer = () => this.notifyListeners(); + this.awareness.on('change', observer); return () => { this.listeners.delete(callback); - this.state.cursors.unobserve(cursorObserver); - this.state.selections.unobserve(selectionObserver); + this.awareness.off('change', observer); }; } + /** Destroy the awareness instance */ + destroy(): void { + this.awareness.destroy(); + } + private notifyListeners(): void { for (const cb of this.listeners) { cb(); diff --git a/frontend/src/crdt/WebSocketProvider.ts b/frontend/src/crdt/WebSocketProvider.ts index 2919eab..717c2a3 100644 --- a/frontend/src/crdt/WebSocketProvider.ts +++ b/frontend/src/crdt/WebSocketProvider.ts @@ -1,9 +1,14 @@ /** * WebSocketProvider – Custom WebSocket provider matching the backend raw-update protocol. * Backend sends Y.encodeStateAsUpdate on connect and applies raw updates on message. + * + * Issue #20: On connect, the client must NOT send its local state immediately. + * Instead, the server sends its state first (Y.encodeStateAsUpdate). The client + * applies that update to sync, then sends only incremental local updates. */ import * as Y from 'yjs'; import type { YjsDocument } from './YjsDocument'; +import type { AwarenessManager } from './AwarenessManager'; export type ConnectionStatus = 'disconnected' | 'connecting' | 'connected' | 'error'; @@ -11,15 +16,22 @@ export interface WebSocketProviderOptions { url: string; docName: string; yjsDoc: YjsDocument; + awareness?: AwarenessManager; onStatusChange?: (status: ConnectionStatus) => void; onSync?: () => void; } +// Message type prefixes for multiplexing Y.Doc updates and awareness updates +// Using a simple byte-prefix scheme: 0 = Y.Doc update, 1 = awareness update +const MSG_DOC_UPDATE = 0; +const MSG_AWARENESS = 1; + export class WebSocketProvider { private ws: WebSocket | null = null; private url: string; private docName: string; private yjsDoc: YjsDocument; + private awareness?: AwarenessManager; private status: ConnectionStatus = 'disconnected'; private onStatusChange?: (status: ConnectionStatus) => void; private onSync?: () => void; @@ -28,11 +40,13 @@ export class WebSocketProvider { private maxReconnectDelay = 30000; private shouldReconnect = true; private authToken: string | undefined; + private synced = false; constructor(opts: WebSocketProviderOptions) { this.url = opts.url; this.docName = opts.docName; this.yjsDoc = opts.yjsDoc; + this.awareness = opts.awareness; this.onStatusChange = opts.onStatusChange; this.onSync = opts.onSync; } @@ -44,6 +58,7 @@ export class WebSocketProvider { this.authToken = token; this.shouldReconnect = true; + this.synced = false; this.setStatus('connecting'); const tokenParam = token ? `?token=${encodeURIComponent(token)}` : ''; @@ -61,20 +76,50 @@ export class WebSocketProvider { this.ws.onopen = () => { this.setStatus('connected'); this.reconnectDelay = 1000; - // Send local state to server for initial sync - const stateUpdate = Y.encodeStateAsUpdate(this.yjsDoc.doc); - if (stateUpdate.length > 0) { - this.ws?.send(stateUpdate); + // Issue #20: Do NOT send local state on connect. + // The server sends its state (Y.encodeStateAsUpdate) to new clients. + // We wait for that server state, apply it, then mark as synced. + // Only after sync do we forward incremental local updates. + // + // If we have awareness, send our local awareness state after connect. + if (this.awareness) { + const awarenessUpdate = this.awareness.encodeLocalState(); + if (awarenessUpdate && awarenessUpdate.length > 0) { + this.sendAwarenessUpdate(awarenessUpdate); + } } - this.onSync?.(); }; this.ws.onmessage = (event: MessageEvent) => { try { const data = new Uint8Array(event.data as ArrayBuffer); - this.yjsDoc.applyUpdate(data); + if (data.length === 0) return; + + // Check message type prefix (first byte) + const msgType = data[0]; + const payload = data.slice(1); + + if (msgType === MSG_AWARENESS) { + // Awareness update — apply to awareness manager + if (this.awareness) { + this.awareness.applyRemoteUpdate(payload, 'remote'); + } + } else { + // Y.Doc update — apply to Y.Doc + // Support both prefixed and non-prefixed messages for backward compatibility + // If the first byte looks like a valid Y.Doc update, use the full data; + // otherwise treat the whole message as the Y.Doc update (legacy mode). + const updateData = msgType === MSG_DOC_UPDATE ? payload : data; + this.yjsDoc.applyUpdate(updateData); + + // Mark as synced after receiving the first server state + if (!this.synced) { + this.synced = true; + this.onSync?.(); + } + } } catch (err) { - console.warn('[WebSocketProvider] Failed to apply update:', err); + console.warn('[WebSocketProvider] Failed to process message:', err); } }; @@ -84,6 +129,7 @@ export class WebSocketProvider { this.ws.onclose = () => { this.ws = null; + this.synced = false; this.setStatus('disconnected'); if (this.shouldReconnect) { this.scheduleReconnect(); @@ -93,8 +139,22 @@ export class WebSocketProvider { /** Send a local Y.Doc update to the server */ sendUpdate(update: Uint8Array): void { + if (this.ws && this.ws.readyState === WebSocket.OPEN && this.synced) { + // Prefix with MSG_DOC_UPDATE byte + const msg = new Uint8Array(update.length + 1); + msg[0] = MSG_DOC_UPDATE; + msg.set(update, 1); + this.ws.send(msg); + } + } + + /** Send a local awareness update to the server */ + sendAwarenessUpdate(update: Uint8Array): void { if (this.ws && this.ws.readyState === WebSocket.OPEN) { - this.ws.send(update); + const msg = new Uint8Array(update.length + 1); + msg[0] = MSG_AWARENESS; + msg.set(update, 1); + this.ws.send(msg); } } @@ -112,6 +172,25 @@ export class WebSocketProvider { }; } + /** Listen for local awareness changes and forward to server */ + bindAwarenessUpdates(): () => void { + if (!this.awareness) return () => {}; + const handler = () => { + const update = this.awareness!.encodeLocalState(); + if (update && update.length > 0) { + this.sendAwarenessUpdate(update); + } + }; + this.awareness.awareness.on('update', handler); + return () => { + this.awareness!.awareness.off('update', handler); + }; + } + + isSynced(): boolean { + return this.synced; + } + getStatus(): ConnectionStatus { return this.status; } @@ -126,6 +205,7 @@ export class WebSocketProvider { this.ws.close(); this.ws = null; } + this.synced = false; this.setStatus('disconnected'); } diff --git a/frontend/src/crdt/index.ts b/frontend/src/crdt/index.ts index 8c77da0..b6b14e5 100644 --- a/frontend/src/crdt/index.ts +++ b/frontend/src/crdt/index.ts @@ -3,5 +3,5 @@ */ export { YjsDocument, type YjsDocumentData } from './YjsDocument'; export { WebSocketProvider, type ConnectionStatus, type WebSocketProviderOptions } from './WebSocketProvider'; -export { AwarenessManager, type UserCursor, type UserSelection, type AwarenessState } from './AwarenessManager'; +export { AwarenessManager, type UserCursor, type UserSelection } from './AwarenessManager'; export { useYjsBinding, type UseYjsBindingOptions, type UseYjsBindingResult } from './useYjsBinding'; diff --git a/frontend/src/crdt/useYjsBinding.ts b/frontend/src/crdt/useYjsBinding.ts index cb5c2ad..fb10288 100644 --- a/frontend/src/crdt/useYjsBinding.ts +++ b/frontend/src/crdt/useYjsBinding.ts @@ -82,17 +82,18 @@ export function useYjsBinding(opts: UseYjsBindingOptions): UseYjsBindingResult { const doc = new YjsDocument(); yjsDocRef.current = doc; + const awareness = new AwarenessManager(doc, userId); + awarenessRef.current = awareness; + const provider = new WebSocketProvider({ url: wsUrl, docName, yjsDoc: doc, + awareness, onStatusChange: (s) => setStatus(s), }); providerRef.current = provider; - const awareness = new AwarenessManager(doc, userId); - awarenessRef.current = awareness; - // Sync local state from Y.Doc on changes const syncState = () => { setElements(doc.getElements()); @@ -123,6 +124,8 @@ export function useYjsBinding(opts: UseYjsBindingOptions): UseYjsBindingResult { // Forward local updates to server const unbindLocal = provider.bindLocalUpdates(); + // Forward local awareness updates to server + const unbindAwarenessUpdates = provider.bindAwarenessUpdates(); // Connect provider.connect(token); @@ -135,7 +138,9 @@ export function useYjsBinding(opts: UseYjsBindingOptions): UseYjsBindingResult { unbindDoc(); unbindAwareness(); unbindLocal(); + unbindAwarenessUpdates(); awareness.removeSelf(); + awareness.destroy(); provider.disconnect(); doc.destroy(); yjsDocRef.current = null;