From 490fb769df1a85a9abc9c2a9ab3396c636c97606 Mon Sep 17 00:00:00 2001 From: Leopoldadmin Date: Mon, 22 Jun 2026 05:52:58 +0000 Subject: [PATCH] feat: Add Yjs WebSocket server implementation --- backend/src/websocket/yjsServer.ts | 107 +++++++++++++++++++++++++++++ 1 file changed, 107 insertions(+) create mode 100644 backend/src/websocket/yjsServer.ts diff --git a/backend/src/websocket/yjsServer.ts b/backend/src/websocket/yjsServer.ts new file mode 100644 index 0000000..4f0f2d9 --- /dev/null +++ b/backend/src/websocket/yjsServer.ts @@ -0,0 +1,107 @@ +import { WebsocketProvider } from 'y-websocket' +import * as Y from 'yjs' +import { LeveldbPersistence } from 'y-leveldb' +import * as fs from 'fs' +import * as path from 'path' + +// Ensure the data directory exists +const yjsDataDir = '/data/yjs' +if (!fs.existsSync(yjsDataDir)) { + fs.mkdirSync(yjsDataDir, { recursive: true }) +} + +// Map to store documents and their persistence +const documents = new Map() + +// Function to get or create a document for a project +async function getDocument(projectId: string): Promise { + const roomName = `project:${projectId}` + + if (documents.has(roomName)) { + return documents.get(roomName)!.doc + } + + // Create a new document + const doc = new Y.Doc() + + // Set up persistence + const dbPath = path.join(yjsDataDir, `${projectId}.ldb`) + const persistence = new LeveldbPersistence(dbPath, doc) + + // Store the document and its persistence + documents.set(roomName, { doc, persistence }) + + return doc +} + +// Function to handle WebSocket connections +async function handleWebSocketConnection(socket: any, request: any) { + try { + // Extract JWT token from query param or header + const token = request.query.token || request.headers.authorization?.replace('Bearer ', '') + + // Verify JWT token (implementation depends on your auth system) + // const user = verifyJWT(token) + + // For now, we'll use a mock user object + const user = { + id: 'user123', + role: 'editor' // or 'viewer' for read-only users + } + + // Extract project ID from the room name + const url = new URL(socket.url, 'ws://localhost') + const roomName = url.searchParams.get('room') + + if (!roomName || !roomName.startsWith('project:')) { + socket.close(4000, 'Invalid room name') + return + } + + const projectId = roomName.replace('project:', '') + + // Get or create the document for this project + const doc = await getDocument(projectId) + + // Set up the WebSocket provider + const provider = new WebsocketProvider('ws://localhost:3001', roomName, doc, { + WebSocket: socket.constructor, + resyncInterval: false, + maxBackoffTime: 2500, + disableBc: true // Disable broadcast channel to ensure all messages go through the server + }) + + // Handle read-only users + if (user.role === 'viewer') { + // Override the send method to prevent updates from read-only users + const originalSend = provider.send + provider.send = function (message: any) { + // Allow sync-step1, sync-step2, and awareness messages + // Block update messages + if (message[0] === 1) { // Update message type + console.log(`Blocking update from read-only user: ${user.id}`) + return + } + originalSend.call(this, message) + } + } + + // Handle awareness updates + provider.awareness.on('update', (changes: any) => { + // Broadcast awareness changes to all clients in the room + provider.broadcastAwareness() + }) + + // Clean up when the connection closes + socket.on('close', () => { + provider.destroy() + }) + + } catch (error) { + console.error('WebSocket connection error:', error) + socket.close(1011, 'Internal server error') + } +} + +// Export functions for use in the Fastify server +export { handleWebSocketConnection, getDocument }