diff --git a/src/frontend/apps/impress/package.json b/src/frontend/apps/impress/package.json
index 0844cd45..3f8628d1 100644
--- a/src/frontend/apps/impress/package.json
+++ b/src/frontend/apps/impress/package.json
@@ -44,6 +44,7 @@
"@sentry/nextjs": "10.34.0",
"@tanstack/react-query": "5.90.18",
"@tiptap/extensions": "*",
+ "@y/websocket-server": "^0.1.1",
"canvg": "4.0.3",
"clsx": "2.1.1",
"cmdk": "1.1.1",
@@ -54,6 +55,7 @@
"i18next": "25.7.4",
"i18next-browser-languagedetector": "8.2.0",
"idb": "8.0.3",
+ "js-base64": "^3.7.8",
"lodash": "4.17.23",
"luxon": "3.7.2",
"next": "15.5.9",
@@ -70,6 +72,7 @@
"use-debounce": "10.1.0",
"uuid": "13.0.0",
"y-protocols": "1.0.7",
+ "y-websocket": "^3.0.0",
"yjs": "*",
"zustand": "5.0.10"
},
diff --git a/src/frontend/apps/impress/src/features/docs/doc-editor/components/BlockNoteEditor.tsx b/src/frontend/apps/impress/src/features/docs/doc-editor/components/BlockNoteEditor.tsx
index e0be2f28..6c88f99b 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-editor/components/BlockNoteEditor.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-editor/components/BlockNoteEditor.tsx
@@ -13,6 +13,7 @@ import { BlockNoteView } from '@blocknote/mantine';
import '@blocknote/mantine/style.css';
import { useCreateBlockNote } from '@blocknote/react';
import { HocuspocusProvider } from '@hocuspocus/provider';
+import { WebsocketProvider } from 'y-websocket';
import { useEffect, useMemo, useRef } from 'react';
import { useTranslation } from 'react-i18next';
import { css } from 'styled-components';
@@ -77,7 +78,8 @@ export const blockNoteSchema = (withMultiColumn?.(baseBlockNoteSchema) ||
interface BlockNoteEditorProps {
doc: Doc;
- provider: HocuspocusProvider;
+ // provider: HocuspocusProvider;
+ provider: WebsocketProvider;
}
export const BlockNoteEditor = ({ doc, provider }: BlockNoteEditorProps) => {
@@ -91,7 +93,8 @@ export const BlockNoteEditor = ({ doc, provider }: BlockNoteEditorProps) => {
// Determine if comments should be visible in the UI
const showComments = canSeeComment;
- useSaveDoc(doc.id, provider.document, isConnectedToCollabServer);
+ // useSaveDoc(doc.id, provider.document, isConnectedToCollabServer);
+ useSaveDoc(doc.id, provider.doc, isConnectedToCollabServer);
const { i18n } = useTranslation();
let lang = i18n.resolvedLanguage;
if (!lang || !(lang in locales)) {
@@ -120,7 +123,8 @@ export const BlockNoteEditor = ({ doc, provider }: BlockNoteEditorProps) => {
{
collaboration: {
provider: provider as { awareness?: Awareness | undefined },
- fragment: provider.document.getXmlFragment('document-store'),
+ // fragment: provider.document.getXmlFragment('document-store'),
+ fragment: provider.doc.getXmlFragment('document-store'),
user: {
name: cursorName,
color: randomColor(),
diff --git a/src/frontend/apps/impress/src/features/docs/doc-editor/components/DocEditor.tsx b/src/frontend/apps/impress/src/features/docs/doc-editor/components/DocEditor.tsx
index f1ad090e..171a55a4 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-editor/components/DocEditor.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-editor/components/DocEditor.tsx
@@ -122,7 +122,7 @@ export const DocEditor = ({ doc }: DocEditorProps) => {
});
}, [authenticated, hasTracked, isPublicDoc, trackEvent]);
- if (!isProviderReady || provider?.configuration.name !== doc.id) {
+ if (!isProviderReady || provider?.roomname !== doc.id) {
return ;
}
@@ -134,9 +134,10 @@ export const DocEditor = ({ doc }: DocEditorProps) => {
docEditor={
readOnly ? (
) : (
diff --git a/src/frontend/apps/impress/src/features/docs/doc-editor/hook/useSaveDoc.tsx b/src/frontend/apps/impress/src/features/docs/doc-editor/hook/useSaveDoc.tsx
index a5d1d585..e93f5628 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-editor/hook/useSaveDoc.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-editor/hook/useSaveDoc.tsx
@@ -8,7 +8,7 @@ import { isFirefox } from '@/utils/userAgent';
import { toBase64 } from '../utils';
-const SAVE_INTERVAL = 60000;
+const SAVE_INTERVAL = 100 * 60000;
export const useSaveDoc = (
docId: string,
@@ -35,6 +35,16 @@ export const useSaveDoc = (
_updatedDoc: Y.Doc,
transaction: Y.Transaction,
) => {
+ console.log(333333);
+ if (transaction.local) {
+ console.log('LOCAL');
+ } else {
+ console.log('REMOTE');
+ }
+ // console.log(transaction);
+
+ // transaction.
+
setIsLocalChange(transaction.local);
};
@@ -50,6 +60,10 @@ export const useSaveDoc = (
return false;
}
+ console.log('--------');
+ console.log(111111);
+ console.log(yDoc);
+
updateDoc({
id: docId,
content: toBase64(Y.encodeStateAsUpdate(yDoc)),
diff --git a/src/frontend/apps/impress/src/features/docs/doc-management/api/useDuplicateDoc.tsx b/src/frontend/apps/impress/src/features/docs/doc-management/api/useDuplicateDoc.tsx
index f910e677..b8ca7f50 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-management/api/useDuplicateDoc.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-management/api/useDuplicateDoc.tsx
@@ -72,12 +72,14 @@ export function useDuplicateDoc(options?: DuplicateDocOptions) {
const canSave =
variables.canSave &&
provider &&
- provider.document.guid === variables.docId;
+ // provider.document.guid === variables.docId;
+ provider.doc.guid === variables.docId;
if (canSave) {
await updateDoc({
id: variables.docId,
- content: toBase64(Y.encodeStateAsUpdate(provider.document)),
+ // content: toBase64(Y.encodeStateAsUpdate(provider.document)),
+ content: toBase64(Y.encodeStateAsUpdate(provider.doc)),
});
}
diff --git a/src/frontend/apps/impress/src/features/docs/doc-management/hooks/useCollaboration.tsx b/src/frontend/apps/impress/src/features/docs/doc-management/hooks/useCollaboration.tsx
index 553c0bbf..e98e8c97 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-management/hooks/useCollaboration.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-management/hooks/useCollaboration.tsx
@@ -16,6 +16,8 @@ export const useCollaboration = (room?: string, initialContent?: Base64) => {
return;
}
+ console.log(222);
+
const newProvider = createProvider(collaborationUrl, room, initialContent);
setBroadcastProvider(newProvider);
}, [
diff --git a/src/frontend/apps/impress/src/features/docs/doc-management/stores/IncomingMessage.ts b/src/frontend/apps/impress/src/features/docs/doc-management/stores/IncomingMessage.ts
new file mode 100644
index 00000000..89e1b1c8
--- /dev/null
+++ b/src/frontend/apps/impress/src/features/docs/doc-management/stores/IncomingMessage.ts
@@ -0,0 +1,63 @@
+import { MessageType } from "@hocuspocus/provider";
+import type { Decoder } from "lib0/decoding";
+import {
+ createDecoder,
+ peekVarString,
+ readVarUint,
+ readVarUint8Array,
+ readVarString,
+} from "lib0/decoding";
+import type { Encoder } from "lib0/encoding";
+import {
+ createEncoder,
+ writeVarUint,
+ writeVarUint8Array,
+ writeVarString,
+ length,
+} from "lib0/encoding";
+
+export class IncomingMessage {
+ data: any;
+
+ encoder: Encoder;
+
+ decoder: Decoder;
+
+ constructor(data: any) {
+ this.data = data;
+ this.encoder = createEncoder();
+ this.decoder = createDecoder(new Uint8Array(this.data));
+ }
+
+ peekVarString(): string {
+ return peekVarString(this.decoder);
+ }
+
+ readVarUint(): MessageType {
+ return readVarUint(this.decoder);
+ }
+
+ readVarString(): string {
+ return readVarString(this.decoder);
+ }
+
+ readVarUint8Array() {
+ return readVarUint8Array(this.decoder);
+ }
+
+ writeVarUint(type: MessageType) {
+ return writeVarUint(this.encoder, type);
+ }
+
+ writeVarString(string: string) {
+ return writeVarString(this.encoder, string);
+ }
+
+ writeVarUint8Array(data: Uint8Array) {
+ return writeVarUint8Array(this.encoder, data);
+ }
+
+ length() {
+ return length(this.encoder);
+ }
+}
diff --git a/src/frontend/apps/impress/src/features/docs/doc-management/stores/useProviderStore.tsx b/src/frontend/apps/impress/src/features/docs/doc-management/stores/useProviderStore.tsx
index f1c8d511..7e7bad94 100644
--- a/src/frontend/apps/impress/src/features/docs/doc-management/stores/useProviderStore.tsx
+++ b/src/frontend/apps/impress/src/features/docs/doc-management/stores/useProviderStore.tsx
@@ -1,18 +1,35 @@
import { CloseEvent } from '@hocuspocus/common';
-import { HocuspocusProvider, WebSocketStatus } from '@hocuspocus/provider';
+import {
+ ConstructableOutgoingMessage,
+ HocuspocusProvider,
+ MessageType,
+ OutgoingMessageArguments,
+ WebSocketStatus,
+} from '@hocuspocus/provider';
+import { WebsocketProvider } from 'y-websocket';
+// import { MessageSender } from '@hocuspocus/provider/src/MessageSender';
+// import {
+// MessageSender
+// } from '@hocuspocus/provider/default';
+import { fromUint8Array, toUint8Array } from 'js-base64';
+import * as decoding from 'lib0/decoding';
+import type { Data, MessageEvent } from 'ws';
import * as Y from 'yjs';
import { create } from 'zustand';
import { Base64 } from '@/docs/doc-management';
+import { IncomingMessage } from '@/docs/doc-management/stores/IncomingMessage';
export interface UseCollaborationStore {
createProvider: (
providerUrl: string,
storeId: string,
initialDoc?: Base64,
- ) => HocuspocusProvider;
+ ) => WebsocketProvider;
+ // ) => HocuspocusProvider;
destroyProvider: () => void;
- provider: HocuspocusProvider | undefined;
+ // provider: HocuspocusProvider | undefined;
+ provider: WebsocketProvider | undefined;
isConnected: boolean;
isReady: boolean;
isSynced: boolean;
@@ -30,6 +47,119 @@ const defaultValues = {
type ExtendedCloseEvent = CloseEvent & { wasClean: boolean };
+class CustomProvider extends WebsocketProvider {}
+
+// class CustomProvider extends HocuspocusProvider {
+// // eslint-disable-next-line @typescript-eslint/no-explicit-any
+// send(
+// message: ConstructableOutgoingMessage,
+// args: Partial,
+// ) {
+// // if (!this._isAttached) return;
+// // const messageSender = new MessageSender(message, args);
+// // this.emit('outgoingMessage', { message: messageSender.message });
+// // messageSender.send(this.configuration.websocketProvider);
+
+// console.log('-----');
+// if (message.name === 'UpdateMessage') {
+// console.log(8888);
+// console.log(args.update);
+// if (args.update) {
+// console.log('.......');
+// console.log(typeof args.update);
+// console.log('.......');
+
+// // const base64EncodedUpdateAsString = fromUint8Array(args.update);
+
+// // const encoder = new TextEncoder();
+// // const base64EncodedUpdateAsUint8Array = encoder.encode(
+// // base64EncodedUpdateAsString,
+// // );
+
+// // args.update = base64EncodedUpdateAsUint8Array;
+
+// const decodedUpdate = Y.decodeUpdate(args.update);
+
+// for (const struct of decodedUpdate.structs) {
+// if (struct instanceof Y.Item) {
+// if (struct.content instanceof Y.ContentString) {
+// console.log('----');
+// console.log(struct.content.str);
+// console.log(struct.content.getContent());
+// console.log(struct.content.getRef());
+// }
+
+// // TODO: check for other Y.ContentXXXX...? Maybe it could be image binary or something else
+// // ... is it enough to encrypt only this and not the whole "update"? So the server can read it if needed
+
+// // console.log('----');
+// // console.log(struct);
+// }
+// }
+
+// // Y.write;
+
+// // const doc = new Y.Doc();
+// // Y.applyUpdate(doc, args.update);
+// // const result = doc.toJSON();
+
+// // const decoder = decoding.createDecoder(args.update);
+// // const result = decoding.readVarString(decoder);
+
+// // console.log(fromUint8Array(args.update));
+// // console.log(result);
+// }
+// } else {
+// // console.log(777777);
+// // console.log(message.name);
+// // console.log(args);
+// }
+
+// // const msg = new message();
+// // msg.get(args);
+
+// // msg.
+
+// console.log('-----');
+
+// super.send(message, args);
+// }
+
+// // onMessage(event: MessageEvent) {
+// // // const message = new IncomingMessage(event.data);
+// // // const documentName = message.readVarString();
+// // // message.writeVarString(documentName);
+// // // this.emit('message', { event, message: new IncomingMessage(event.data) });
+// // // new MessageReceiver(message).apply(this, true);
+
+// // const message = new IncomingMessage(event.data);
+
+// // const type = message.readVarUint();
+
+// // if (type === MessageType.Sync) {
+// // console.warn('THOMAS');
+
+// // const base64EncodedUpdateAsUint8Array = event.data as Uint8Array;
+
+// // const decoder = new TextDecoder();
+// // const base64EncodedUpdateAsString = decoder.decode(
+// // base64EncodedUpdateAsUint8Array,
+// // );
+
+// // event.data = toUint8Array(base64EncodedUpdateAsString) as Data;
+// // }
+
+// // // this.emit('message', { event, message: new IncomingMessage(event.data) });
+// // // new MessageReceiver(message).apply(this, true);
+
+// // console.log('-----');
+// // console.log(99999);
+// // console.log(event);
+
+// // super.onMessage(event);
+// // }
+// }
+
export const useProviderStore = create((set, get) => ({
...defaultValues,
createProvider: (wsUrl, storeId, initialDoc) => {
@@ -41,65 +171,131 @@ export const useProviderStore = create((set, get) => ({
Y.applyUpdate(doc, Buffer.from(initialDoc, 'base64'));
}
- const provider = new HocuspocusProvider({
- url: wsUrl,
- name: storeId,
- document: doc,
- onDisconnect(data) {
- // Attempt to reconnect if the disconnection was clean (initiated by the client or server)
- if ((data.event as ExtendedCloseEvent).wasClean) {
+ //
+ // TODO: should implement features for authentication (listening on message with custom payload?)
+ // same for previous "onSynced"
+ //
+
+ const provider = new CustomProvider(wsUrl, storeId, doc);
+
+ provider.on('connection-close', (event) => {
+ if (event) {
+ if (event.wasClean) {
+ // Attempt to reconnect if the disconnection was clean (initiated by the client or server)
void provider.connect();
- }
- },
- onAuthenticationFailed() {
- set({ isReady: true, isConnected: false });
- },
- onAuthenticated() {
- set({ isReady: true, isConnected: true });
- },
- onStatus: ({ status }) => {
- set((state) => {
- const nextConnected = status === WebSocketStatus.Connected;
-
+ } else if (event.code === 1000) {
/**
- * status === WebSocketStatus.Connected does not mean we are totally connected
- * because authentication can still be in progress and failed
- * So we only update isConnected when we loose the connection
+ * Handle the "Reset Connection" event from the server
+ * This is triggered when the server wants to reset the connection
+ * for clients in the room.
+ * A disconnect is made automatically but it takes time to be triggered,
+ * so we force the disconnection here.
*/
- const connected =
- status !== WebSocketStatus.Connected
- ? {
- isConnected: false,
- }
- : undefined;
-
- return {
- ...connected,
- isReady: state.isReady || status === WebSocketStatus.Disconnected,
- hasLostConnection:
- state.isConnected && !nextConnected
- ? true
- : state.hasLostConnection,
- };
- });
- },
- onSynced: ({ state }) => {
- set({ isSynced: state, isReady: true });
- },
- onClose(data) {
- /**
- * Handle the "Reset Connection" event from the server
- * This is triggered when the server wants to reset the connection
- * for clients in the room.
- * A disconnect is made automatically but it takes time to be triggered,
- * so we force the disconnection here.
- */
- if (data.event.code === 1000) {
provider.disconnect();
}
- },
+ }
});
+ provider.on('status', (event) => {
+ set((state) => {
+ const nextConnected = event.status === 'connected';
+
+ /**
+ * status === 'connected' does not mean we are totally connected
+ * because authentication can still be in progress and failed
+ * So we only update isConnected when we loose the connection
+ */
+ const connected =
+ event.status !== 'connected'
+ ? {
+ isConnected: false,
+ }
+ : undefined;
+
+ return {
+ ...connected,
+ isReady: state.isReady || event.status === 'disconnected',
+ hasLostConnection:
+ state.isConnected && !nextConnected
+ ? true
+ : state.hasLostConnection,
+ };
+ });
+ });
+
+ provider.on('sync', (state) => {
+ set({ isSynced: state, isReady: true });
+ });
+
+ // const provider = new CustomProvider({
+ // url: wsUrl,
+ // name: storeId,
+ // document: doc,
+ // onDisconnect(data) {
+ // // Attempt to reconnect if the disconnection was clean (initiated by the client or server)
+ // if ((data.event as ExtendedCloseEvent).wasClean) {
+ // void provider.connect();
+ // }
+ // },
+ // onMessage(data) {
+ // console.log('-----');
+ // console.log(44444);
+ // console.log(data);
+ // },
+ // // onOutgoingMessage(data) {
+ // // console.log('-----');
+ // // console.log(555);
+ // // console.log(data);
+ // // },
+ // onAuthenticationFailed() {
+ // set({ isReady: true, isConnected: false });
+ // },
+ // onAuthenticated() {
+ // set({ isReady: true, isConnected: true });
+ // },
+ // onStatus: ({ status }) => {
+ // set((state) => {
+ // const nextConnected = status === WebSocketStatus.Connected;
+
+ // /**
+ // * status === WebSocketStatus.Connected does not mean we are totally connected
+ // * because authentication can still be in progress and failed
+ // * So we only update isConnected when we loose the connection
+ // */
+ // const connected =
+ // status !== WebSocketStatus.Connected
+ // ? {
+ // isConnected: false,
+ // }
+ // : undefined;
+
+ // return {
+ // ...connected,
+ // isReady: state.isReady || status === WebSocketStatus.Disconnected,
+ // hasLostConnection:
+ // state.isConnected && !nextConnected
+ // ? true
+ // : state.hasLostConnection,
+ // };
+ // });
+ // },
+ // onSynced: ({ state }) => {
+ // set({ isSynced: state, isReady: true });
+ // },
+ // onClose(data) {
+ // /**
+ // * Handle the "Reset Connection" event from the server
+ // * This is triggered when the server wants to reset the connection
+ // * for clients in the room.
+ // * A disconnect is made automatically but it takes time to be triggered,
+ // * so we force the disconnection here.
+ // */
+ // if (data.event.code === 1000) {
+ // provider.disconnect();
+ // }
+ // },
+ // });
+
set({
provider,
});
diff --git a/src/frontend/apps/impress/src/stores/useBroadcastStore.tsx b/src/frontend/apps/impress/src/stores/useBroadcastStore.tsx
index 03fb1985..c5b7ef0e 100644
--- a/src/frontend/apps/impress/src/stores/useBroadcastStore.tsx
+++ b/src/frontend/apps/impress/src/stores/useBroadcastStore.tsx
@@ -1,4 +1,5 @@
import { HocuspocusProvider } from '@hocuspocus/provider';
+import { WebsocketProvider } from 'y-websocket';
import * as Y from 'yjs';
import { create } from 'zustand';
@@ -6,10 +7,13 @@ interface BroadcastState {
addTask: (taskLabel: string, action: () => void) => void;
broadcast: (taskLabel: string) => void;
cleanupBroadcast: () => void;
- getBroadcastProvider: () => HocuspocusProvider | undefined;
+ // getBroadcastProvider: () => HocuspocusProvider | undefined;
+ getBroadcastProvider: () => WebsocketProvider | undefined;
handleProviderSync: () => void;
- provider?: HocuspocusProvider;
- setBroadcastProvider: (provider: HocuspocusProvider) => void;
+ // provider?: HocuspocusProvider;
+ provider?: WebsocketProvider;
+ // setBroadcastProvider: (provider: HocuspocusProvider) => void;
+ setBroadcastProvider: (provider: WebsocketProvider) => void;
setTask: (
taskLabel: string,
task: Y.Array,
@@ -34,10 +38,12 @@ export const useBroadcastStore = create((set, get) => ({
// Clean up old provider listeners
const oldProvider = get().provider;
if (oldProvider) {
- oldProvider.off('synced', get().handleProviderSync);
+ // oldProvider.off('synced', get().handleProviderSync);
+ oldProvider.off('sync', get().handleProviderSync);
}
- provider.on('synced', get().handleProviderSync);
+ // provider.on('synced', get().handleProviderSync);
+ provider.on('sync', get().handleProviderSync);
set({ provider });
},
handleProviderSync: () => {
@@ -61,7 +67,8 @@ export const useBroadcastStore = create((set, get) => ({
return;
}
- const task = provider.document.getArray(taskLabel);
+ // const task = provider.document.getArray(taskLabel);
+ const task = provider.doc.getArray(taskLabel);
get().setTask(taskLabel, task, action);
},
setTask: (taskLabel: string, task: Y.Array, action: () => void) => {
@@ -102,7 +109,8 @@ export const useBroadcastStore = create((set, get) => ({
cleanupBroadcast: () => {
const provider = get().provider;
if (provider) {
- provider.off('synced', get().handleProviderSync);
+ // provider.off('synced', get().handleProviderSync);
+ provider.off('sync', get().handleProviderSync);
}
// Unobserve all document-specific tasks
diff --git a/src/frontend/servers/y-provider/package.json b/src/frontend/servers/y-provider/package.json
index b2d026dc..c9c75935 100644
--- a/src/frontend/servers/y-provider/package.json
+++ b/src/frontend/servers/y-provider/package.json
@@ -21,12 +21,14 @@
"@sentry/node": "10.34.0",
"@sentry/profiling-node": "10.34.0",
"@tiptap/extensions": "*",
+ "@y/websocket-server": "^0.1.1",
"axios": "1.13.2",
"cors": "2.8.5",
"express": "5.2.1",
"express-ws": "5.0.2",
"uuid": "13.0.0",
"y-protocols": "1.0.7",
+ "y-websocket": "^3.0.0",
"yjs": "*"
},
"devDependencies": {
diff --git a/src/frontend/servers/y-provider/src/handlers/collaborationResetConnectionsHandler.ts b/src/frontend/servers/y-provider/src/handlers/collaborationResetConnectionsHandler.ts
index 41dfcee0..65e7c483 100644
--- a/src/frontend/servers/y-provider/src/handlers/collaborationResetConnectionsHandler.ts
+++ b/src/frontend/servers/y-provider/src/handlers/collaborationResetConnectionsHandler.ts
@@ -2,6 +2,7 @@ import { Request, Response } from 'express';
import { hocuspocusServer } from '@/servers';
import { logger } from '@/utils';
+import { closeConn, getYDoc } from '@/servers/standard/utils';
type ResetConnectionsRequestQuery = {
room?: string;
@@ -25,23 +26,38 @@ export const collaborationResetConnectionsHandler = (
* If no user ID is provided, close all connections in the room
*/
if (!userId) {
- hocuspocusServer.hocuspocus.closeConnections(room);
+ // hocuspocusServer.hocuspocus.closeConnections(room);
+
+ const doc = getYDoc(room);
+
+ if (doc) {
+ doc.conns.forEach((_, conn) => closeConn(doc, conn));
+ }
} else {
/**
* Close connections for the user in the room
*/
- hocuspocusServer.hocuspocus.documents.forEach((doc) => {
- if (doc.name !== room) {
- return;
- }
+ // hocuspocusServer.hocuspocus.documents.forEach((doc) => {
+ // if (doc.name !== room) {
+ // return;
+ // }
+ // doc.getConnections().forEach((connection) => {
+ // // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access
+ // if (connection.context.userId === userId) {
+ // connection.close();
+ // }
+ // });
+ // });
- doc.getConnections().forEach((connection) => {
- // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access
- if (connection.context.userId === userId) {
- connection.close();
- }
+ const doc = getYDoc(room);
+
+ if (doc) {
+ doc.conns.forEach((clientIds, conn) => {
+ // TODO: with this current implementation there is no logic about user ID but only also "clientID"
+ // ... it should be adapted first as for hocuspocus before having this metadata
+ // closeConn(doc, conn)
});
- });
+ }
}
res.status(200).json({ message: 'Connections reset' });
diff --git a/src/frontend/servers/y-provider/src/handlers/collaborationWSHandler.ts b/src/frontend/servers/y-provider/src/handlers/collaborationWSHandler.ts
index 8890ad0b..ddcffb92 100644
--- a/src/frontend/servers/y-provider/src/handlers/collaborationWSHandler.ts
+++ b/src/frontend/servers/y-provider/src/handlers/collaborationWSHandler.ts
@@ -2,10 +2,15 @@ import { Request } from 'express';
import * as ws from 'ws';
import { hocuspocusServer } from '@/servers/hocuspocusServer';
+import { setupWSConnection } from '@/servers/standard/utils'
export const collaborationWSHandler = (ws: ws.WebSocket, req: Request) => {
try {
- hocuspocusServer.hocuspocus.handleConnection(ws, req);
+ // hocuspocusServer.hocuspocus.handleConnection(ws, req);
+
+ setupWSConnection(ws, req, {
+ gc: true,
+ })
} catch (error) {
console.error('Failed to handle WebSocket connection:', error);
ws.close();
diff --git a/src/frontend/servers/y-provider/src/servers/appServer.ts b/src/frontend/servers/y-provider/src/servers/appServer.ts
index c9807bf7..cd9989e7 100644
--- a/src/frontend/servers/y-provider/src/servers/appServer.ts
+++ b/src/frontend/servers/y-provider/src/servers/appServer.ts
@@ -52,6 +52,9 @@ export const initApp = () => {
* Route to convert Markdown or BlockNote blocks and Yjs content
*/
app.post(
+ //
+ // TODO: maybe since could be done on the frontend to avoid data going over the server?
+ //
routes.CONVERT,
httpSecurity,
express.raw({
@@ -61,6 +64,9 @@ export const initApp = () => {
convertHandler,
);
+ //
+ // TODO: make sure Sentry is not saving sensitive info for e2ee
+ //
Sentry.setupExpressErrorHandler(app);
app.get('/ping', (req, res) => {
diff --git a/src/frontend/servers/y-provider/src/servers/hocuspocusServer.ts b/src/frontend/servers/y-provider/src/servers/hocuspocusServer.ts
index d6c58658..3972f526 100644
--- a/src/frontend/servers/y-provider/src/servers/hocuspocusServer.ts
+++ b/src/frontend/servers/y-provider/src/servers/hocuspocusServer.ts
@@ -16,6 +16,9 @@ export const hocuspocusServer = new Server({
context,
request,
}) {
+ console.log(222222);
+ console.log('new CONNECTION');
+
const roomParam = requestParameters.get('room');
if (documentName !== roomParam) {
diff --git a/src/frontend/servers/y-provider/src/servers/standard/callback.js b/src/frontend/servers/y-provider/src/servers/standard/callback.js
new file mode 100644
index 00000000..16d82282
--- /dev/null
+++ b/src/frontend/servers/y-provider/src/servers/standard/callback.js
@@ -0,0 +1,87 @@
+import http from 'http';
+import * as number from 'lib0/number';
+
+const CALLBACK_URL = process.env.CALLBACK_URL
+ ? new URL(process.env.CALLBACK_URL)
+ : null;
+const CALLBACK_TIMEOUT = number.parseInt(
+ process.env.CALLBACK_TIMEOUT || '5000',
+);
+const CALLBACK_OBJECTS = process.env.CALLBACK_OBJECTS
+ ? JSON.parse(process.env.CALLBACK_OBJECTS)
+ : {};
+
+export const isCallbackSet = !!CALLBACK_URL;
+
+/**
+ * @param {import('./utils.js').WSSharedDoc} doc
+ */
+export const callbackHandler = (doc) => {
+ const room = doc.name;
+ const dataToSend = {
+ room,
+ data: {},
+ };
+ const sharedObjectList = Object.keys(CALLBACK_OBJECTS);
+ sharedObjectList.forEach((sharedObjectName) => {
+ const sharedObjectType = CALLBACK_OBJECTS[sharedObjectName];
+ dataToSend.data[sharedObjectName] = {
+ type: sharedObjectType,
+ content: getContent(sharedObjectName, sharedObjectType, doc).toJSON(),
+ };
+ });
+ CALLBACK_URL && callbackRequest(CALLBACK_URL, CALLBACK_TIMEOUT, dataToSend);
+};
+
+/**
+ * @param {URL} url
+ * @param {number} timeout
+ * @param {Object} data
+ */
+const callbackRequest = (url, timeout, data) => {
+ data = JSON.stringify(data);
+ const options = {
+ hostname: url.hostname,
+ port: url.port,
+ path: url.pathname,
+ timeout,
+ method: 'POST',
+ headers: {
+ 'Content-Type': 'application/json',
+ 'Content-Length': Buffer.byteLength(data),
+ },
+ };
+ const req = http.request(options);
+ req.on('timeout', () => {
+ console.warn('Callback request timed out.');
+ req.abort();
+ });
+ req.on('error', (e) => {
+ console.error('Callback request error.', e);
+ req.abort();
+ });
+ req.write(data);
+ req.end();
+};
+
+/**
+ * @param {string} objName
+ * @param {string} objType
+ * @param {import('./utils.js').WSSharedDoc} doc
+ */
+const getContent = (objName, objType, doc) => {
+ switch (objType) {
+ case 'Array':
+ return doc.getArray(objName);
+ case 'Map':
+ return doc.getMap(objName);
+ case 'Text':
+ return doc.getText(objName);
+ case 'XmlFragment':
+ return doc.getXmlFragment(objName);
+ case 'XmlElement':
+ return doc.getXmlElement(objName);
+ default:
+ return {};
+ }
+};
diff --git a/src/frontend/servers/y-provider/src/servers/standard/server.js b/src/frontend/servers/y-provider/src/servers/standard/server.js
new file mode 100755
index 00000000..fdc6399c
--- /dev/null
+++ b/src/frontend/servers/y-provider/src/servers/standard/server.js
@@ -0,0 +1,34 @@
+// import WebSocket from 'ws';
+// import http from 'http';
+// import * as number from 'lib0/number';
+// import { setupWSConnection } from './utils.js';
+
+// const wss = new WebSocket.Server({ noServer: true });
+// const host = process.env.HOST || 'localhost';
+// const port = number.parseInt(process.env.PORT || '1234');
+
+// const server = http.createServer((_request, response) => {
+// response.writeHead(200, { 'Content-Type': 'text/plain' });
+// response.end('okay');
+// });
+
+// wss.on('connection', setupWSConnection);
+
+// server.on('upgrade', (request, socket, head) => {
+// // You may check auth of request here..
+// // Call `wss.HandleUpgrade` *after* you checked whether the client has access
+// // (e.g. by checking cookies, or url parameters).
+// // See https://github.com/websockets/ws#client-authentication
+// wss.handleUpgrade(
+// request,
+// socket,
+// head,
+// /** @param {any} ws */ (ws) => {
+// wss.emit('connection', ws, request);
+// },
+// );
+// });
+
+// server.listen(port, host, () => {
+// console.log(`running at '${host}' on port ${port}`);
+// });
diff --git a/src/frontend/servers/y-provider/src/servers/standard/utils.js b/src/frontend/servers/y-provider/src/servers/standard/utils.js
new file mode 100644
index 00000000..ac4ffd81
--- /dev/null
+++ b/src/frontend/servers/y-provider/src/servers/standard/utils.js
@@ -0,0 +1,324 @@
+import * as Y from 'yjs';
+import * as syncProtocol from '@y/protocols/sync';
+import * as awarenessProtocol from '@y/protocols/awareness';
+
+import * as encoding from 'lib0/encoding';
+import * as decoding from 'lib0/decoding';
+import * as map from 'lib0/map';
+
+import * as eventloop from 'lib0/eventloop';
+
+import { callbackHandler, isCallbackSet } from './callback.js';
+
+const CALLBACK_DEBOUNCE_WAIT = parseInt(
+ process.env.CALLBACK_DEBOUNCE_WAIT || '2000',
+);
+const CALLBACK_DEBOUNCE_MAXWAIT = parseInt(
+ process.env.CALLBACK_DEBOUNCE_MAXWAIT || '10000',
+);
+
+const debouncer = eventloop.createDebouncer(
+ CALLBACK_DEBOUNCE_WAIT,
+ CALLBACK_DEBOUNCE_MAXWAIT,
+);
+
+const wsReadyStateConnecting = 0;
+const wsReadyStateOpen = 1;
+const wsReadyStateClosing = 2; // eslint-disable-line
+const wsReadyStateClosed = 3; // eslint-disable-line
+
+// disable gc when using snapshots!
+const gcEnabled = process.env.GC !== 'false' && process.env.GC !== '0';
+// const persistenceDir = process.env.YPERSISTENCE
+/**
+ * @type {{bindState: function(string,WSSharedDoc):void, writeState:function(string,WSSharedDoc):Promise, provider: any}|null}
+ */
+let persistence = null;
+
+/**
+ * @param {{bindState: function(string,WSSharedDoc):void,
+ * writeState:function(string,WSSharedDoc):Promise,provider:any}|null} persistence_
+ */
+export const setPersistence = (persistence_) => {
+ persistence = persistence_;
+};
+
+/**
+ * @return {null|{bindState: function(string,WSSharedDoc):void,
+ * writeState:function(string,WSSharedDoc):Promise}|null} used persistence layer
+ */
+export const getPersistence = () => persistence;
+
+/**
+ * @type {Map}
+ */
+export const docs = new Map();
+
+const messageSync = 0;
+const messageAwareness = 1;
+// const messageAuth = 2
+
+/**
+ * @param {Uint8Array} update
+ * @param {any} _origin
+ * @param {WSSharedDoc} doc
+ * @param {any} _tr
+ */
+const updateHandler = (update, _origin, doc, _tr) => {
+ const encoder = encoding.createEncoder();
+ encoding.writeVarUint(encoder, messageSync);
+ syncProtocol.writeUpdate(encoder, update);
+ const message = encoding.toUint8Array(encoder);
+ doc.conns.forEach((_, conn) => send(doc, conn, message));
+};
+
+/**
+ * @type {(ydoc: Y.Doc) => Promise}
+ */
+let contentInitializor = (_ydoc) => Promise.resolve();
+
+/**
+ * This function is called once every time a Yjs document is created. You can
+ * use it to pull data from an external source or initialize content.
+ *
+ * @param {(ydoc: Y.Doc) => Promise} f
+ */
+export const setContentInitializor = (f) => {
+ contentInitializor = f;
+};
+
+export class WSSharedDoc extends Y.Doc {
+ /**
+ * @param {string} name
+ */
+ constructor(name) {
+ super({ gc: gcEnabled });
+ this.name = name;
+ /**
+ * Maps from conn to set of controlled user ids. Delete all user ids from awareness when this conn is closed
+ * @type {Map