✨(frontend) create class CollaborationProvider

Create the CollaborationProvider class.
This class is inherited from HocuspocusProvider class.
This class integrate a fallback mechanism to handle the
case where the user cannot connect with websockets.
It will use post request to send the data to the
collaboration server.
It will use an EventSource to receive the data from the
collaboration server.
This commit is contained in:
Anthony LC
2025-02-20 11:20:50 +01:00
parent f27e968c15
commit a2f1e32f21
11 changed files with 704 additions and 70 deletions
+5
View File
@@ -8,6 +8,10 @@ and this project adheres to
## [Unreleased]
## Added
- ✨Collaboration long polling fallback #517
## Changed
- 🛂(frontend) Restore version visibility #629
@@ -18,6 +22,7 @@ and this project adheres to
- ♻️(frontend) improve table pdf rendering
## [2.2.0] - 2025-02-10
## Added
@@ -0,0 +1,97 @@
import { expect, test } from '@playwright/test';
import { createDoc, verifyDocName } from './common';
test.beforeEach(async ({ page }) => {
await page.goto('/');
});
test.describe('Doc Collaboration', () => {
/**
* We check:
* - connection to the collaborative server
* - signal of the backend to the collaborative server (connection should close)
* - reconnection to the collaborative server
*/
test('checks the connection with collaborative server', async ({
page,
browserName,
}) => {
let webSocketPromise = page.waitForEvent('websocket', (webSocket) => {
return webSocket
.url()
.includes('ws://localhost:8083/collaboration/ws/?room=');
});
const [title] = await createDoc(page, 'doc-editor', browserName, 1);
await verifyDocName(page, title);
let webSocket = await webSocketPromise;
expect(webSocket.url()).toContain(
'ws://localhost:8083/collaboration/ws/?room=',
);
// Is connected
let framesentPromise = webSocket.waitForEvent('framesent');
await page.locator('.ProseMirror.bn-editor').click();
await page.locator('.ProseMirror.bn-editor').fill('Hello World');
let framesent = await framesentPromise;
expect(framesent.payload).not.toBeNull();
await page.getByRole('button', { name: 'Share' }).click();
const selectVisibility = page.getByLabel('Visibility', { exact: true });
// When the visibility is changed, the ws should closed the connection (backend signal)
const wsClosePromise = webSocket.waitForEvent('close');
await selectVisibility.click();
await page
.getByRole('button', {
name: 'Connected',
})
.click();
// Assert that the doc reconnects to the ws
const wsClose = await wsClosePromise;
expect(wsClose.isClosed()).toBeTruthy();
// Checkt the ws is connected again
webSocketPromise = page.waitForEvent('websocket', (webSocket) => {
return webSocket
.url()
.includes('ws://localhost:8083/collaboration/ws/?room=');
});
webSocket = await webSocketPromise;
framesentPromise = webSocket.waitForEvent('framesent');
framesent = await framesentPromise;
expect(framesent.payload).not.toBeNull();
});
test('checks the connection switch to polling after websocket failure', async ({
page,
browserName,
}) => {
const responsePromise = page.waitForResponse(
(response) =>
response.url().includes('/poll/') && response.status() === 200,
);
await page.routeWebSocket(
'ws://localhost:8083/collaboration/ws/**',
async (ws) => {
await ws.close();
},
);
await page.reload();
await createDoc(page, 'doc-polling', browserName, 1);
const response = await responsePromise;
expect(response.ok()).toBeTruthy();
});
});
@@ -88,70 +88,6 @@ test.describe('Doc Editor', () => {
).toBeVisible();
});
/**
* We check:
* - connection to the collaborative server
* - signal of the backend to the collaborative server (connection should close)
* - reconnection to the collaborative server
*/
test('checks the connection with collaborative server', async ({
page,
browserName,
}) => {
let webSocketPromise = page.waitForEvent('websocket', (webSocket) => {
return webSocket
.url()
.includes('ws://localhost:8083/collaboration/ws/?room=');
});
const randomDoc = await createDoc(page, 'doc-editor', browserName, 1);
await verifyDocName(page, randomDoc[0]);
let webSocket = await webSocketPromise;
expect(webSocket.url()).toContain(
'ws://localhost:8083/collaboration/ws/?room=',
);
// Is connected
let framesentPromise = webSocket.waitForEvent('framesent');
await page.locator('.ProseMirror.bn-editor').click();
await page.locator('.ProseMirror.bn-editor').fill('Hello World');
let framesent = await framesentPromise;
expect(framesent.payload).not.toBeNull();
await page.getByRole('button', { name: 'Share' }).click();
const selectVisibility = page.getByLabel('Visibility', { exact: true });
// When the visibility is changed, the ws should closed the connection (backend signal)
const wsClosePromise = webSocket.waitForEvent('close');
await selectVisibility.click();
await page
.getByRole('button', {
name: 'Connected',
})
.click();
// Assert that the doc reconnects to the ws
const wsClose = await wsClosePromise;
expect(wsClose.isClosed()).toBeTruthy();
// Checkt the ws is connected again
webSocketPromise = page.waitForEvent('websocket', (webSocket) => {
return webSocket
.url()
.includes('ws://localhost:8083/collaboration/ws/?room=');
});
webSocket = await webSocketPromise;
framesentPromise = webSocket.waitForEvent('framesent');
framesent = await framesentPromise;
expect(framesent.payload).not.toBeNull();
});
test('markdown button converts from markdown to the editor syntax json', async ({
page,
browserName,
+1
View File
@@ -62,6 +62,7 @@
"@types/node": "*",
"@types/react": "18.3.12",
"@types/react-dom": "*",
"@types/ws": "8.5.13",
"cross-env": "7.0.3",
"dotenv": "16.4.7",
"eslint-config-impress": "*",
@@ -0,0 +1,67 @@
import { APIError, errorCauses } from '@/api';
interface PollOutgoingMessageParams {
pollUrl: string;
message64: string;
}
interface PollOutgoingMessageResponse {
updated?: boolean;
}
export const pollOutgoingMessageRequest = async ({
pollUrl,
message64,
}: PollOutgoingMessageParams): Promise<PollOutgoingMessageResponse> => {
const response = await fetch(pollUrl, {
method: 'POST',
credentials: 'include',
headers: {
'Content-Type': 'application/json',
},
body: JSON.stringify({
message64,
}),
});
if (!response.ok) {
throw new APIError(
`Post poll message request failed`,
await errorCauses(response),
);
}
return response.json() as Promise<PollOutgoingMessageResponse>;
};
interface PollSyncParams {
pollUrl: string;
localDoc64: string;
}
interface PollSyncResponse {
syncDoc64?: string;
}
export const postPollSyncRequest = async ({
pollUrl,
localDoc64,
}: PollSyncParams): Promise<PollSyncResponse> => {
const response = await fetch(pollUrl, {
method: 'POST',
credentials: 'include',
headers: {
'Content-Type': 'application/json',
},
body: JSON.stringify({
localDoc64,
}),
});
if (!response.ok) {
throw new APIError(
`Sync request failed: ${response.status} ${response.statusText}`,
await errorCauses(response),
);
}
return response.json() as Promise<PollSyncResponse>;
};
@@ -6,17 +6,29 @@ import { useBroadcastStore } from '@/stores';
import { useProviderStore } from '../stores/useProviderStore';
import { Base64 } from '../types';
export const useCollaboration = (room?: string, initialContent?: Base64) => {
export const useCollaboration = (
room?: string,
initialContent?: Base64,
canEdit?: boolean,
) => {
const collaborationUrl = useCollaborationUrl(room);
const { setBroadcastProvider } = useBroadcastStore();
const { provider, createProvider, destroyProvider } = useProviderStore();
/**
* Initialize the provider
*/
useEffect(() => {
if (!room || !collaborationUrl || provider) {
if (!room || !collaborationUrl || provider || canEdit === undefined) {
return;
}
const newProvider = createProvider(collaborationUrl, room, initialContent);
const newProvider = createProvider(
collaborationUrl,
room,
canEdit,
initialContent,
);
setBroadcastProvider(newProvider);
}, [
provider,
@@ -25,6 +37,7 @@ export const useCollaboration = (room?: string, initialContent?: Base64) => {
initialContent,
createProvider,
setBroadcastProvider,
canEdit,
]);
/**
@@ -0,0 +1,325 @@
import crypto from 'crypto';
import {
CompleteHocuspocusProviderConfiguration,
CompleteHocuspocusProviderWebsocketConfiguration,
HocuspocusProvider,
HocuspocusProviderConfiguration,
WebSocketStatus,
onOutgoingMessageParameters,
onStatusParameters,
} from '@hocuspocus/provider';
import type { MessageEvent } from 'ws';
import * as Y from 'yjs';
import { isAPIError } from '@/api';
import {
pollOutgoingMessageRequest,
postPollSyncRequest,
} from '../api/collaborationRequests';
import { toBase64 } from '../utils';
type HocuspocusProviderConfigurationUrl = Required<
Pick<CompleteHocuspocusProviderConfiguration, 'name'>
> &
Partial<CompleteHocuspocusProviderConfiguration> &
Required<Pick<CompleteHocuspocusProviderWebsocketConfiguration, 'url'>>;
export const isHocuspocusProviderConfigurationUrl = (
data: HocuspocusProviderConfiguration,
): data is HocuspocusProviderConfigurationUrl => {
return 'url' in data;
};
type CollaborationProviderConfiguration = HocuspocusProviderConfiguration & {
canEdit: boolean;
};
export class CollaborationProvider extends HocuspocusProvider {
/**
* If the user can edit the document
*/
public canEdit = false;
/**
* If the long polling is started
* it is used to avoid starting it multiple times
* when the websocket is failed.
*/
public isLongPollingStarted = false;
/**
* If the document is syncing with the server
* it is used to avoid starting it multiple times.
*/
public isSyncing = false;
/**
* The document can pass out of sync
* then sync again with a next updates so
* we add a counter to avoid syncing the document
* to quickly.
*/
public seemsUnsyncCount = 0;
public seemsUnsyncMaxCount = 5;
/**
* In Safari or Firefox the websocket takes time before passing in
* mode failed, it can takes up to 1 minutes. To avoid this latence
* we set isWebsocketFailed to true, it is connects it will switch
* to false.
*/
public isWebsocketFailed = true;
/**
* There is a ping-pong mechanism with awareness, receipt awareness is send again,
* it creates useless requests.
* We use this variable to avoid treating the same awareness message twice.
*/
private treatedAwarenessMessage: string | null = null;
/**
* Server-Sent Events
*/
protected sse: EventSource | null = null;
/**
* Polling timeout
* It is used to avoid starting the polling to quickly
* to let the class init properly.
*/
protected pollTimeout: NodeJS.Timeout | null = null;
/**
* Easy way to get the url of the server
*/
protected url = '';
public constructor(configuration: CollaborationProviderConfiguration) {
let url = '';
if (isHocuspocusProviderConfigurationUrl(configuration)) {
url = configuration.url;
}
super(configuration);
this.url = url;
this.canEdit = configuration.canEdit;
if (configuration.canEdit) {
this.on('outgoingMessage', this.onPollOutgoingMessage.bind(this));
}
}
public setPollDefaultValues(): void {
this.isLongPollingStarted = false;
this.isWebsocketFailed = false;
this.seemsUnsyncCount = 0;
this.sse?.close();
this.sse = null;
if (this.pollTimeout) {
clearTimeout(this.pollTimeout);
}
}
public destroy(): void {
super.destroy();
this.setPollDefaultValues();
}
public onStatus({ status }: onStatusParameters) {
if (status === WebSocketStatus.Connecting) {
this.isWebsocketFailed = true;
if (this.pollTimeout) {
clearTimeout(this.pollTimeout);
}
this.pollTimeout = setTimeout(() => {
this.initPolling();
}, 5000);
} else if (status === WebSocketStatus.Connected) {
this.setPollDefaultValues();
}
super.onStatus({ status });
}
public initPolling() {
if (this.isLongPollingStarted || !this.isWebsocketFailed) {
return;
}
this.isLongPollingStarted = true;
void this.pollSync(true);
this.initCollaborationSSE();
}
protected toPollUrl(endpoint: string): string {
let pollUrl = this.url.replace('ws:', 'http:');
if (pollUrl.includes('wss:')) {
pollUrl = pollUrl.replace('wss:', 'https:');
}
pollUrl = pollUrl.replace('/ws/', '/ws/poll/' + endpoint + '/');
// To have our requests not cached
return `${pollUrl}&${Date.now()}`;
}
protected isDuplicateAwareness(message64: string): boolean {
if (this.treatedAwarenessMessage === message64) {
return true;
}
this.treatedAwarenessMessage = message64;
return false;
}
/**
* Outgoing message event
*
* Sent to the server the message to
* be sent to the other users
*/
public async onPollOutgoingMessage({ message }: onOutgoingMessageParameters) {
if (!this.isWebsocketFailed || !this.canEdit) {
return;
}
const message64 = Buffer.from(message.toUint8Array()).toString('base64');
if (this.isDuplicateAwareness(message64)) {
return;
}
try {
const { updated } = await pollOutgoingMessageRequest({
pollUrl: this.toPollUrl('message'),
message64,
});
if (!updated) {
await this.pollSync();
}
} catch (error: unknown) {
if (isAPIError(error)) {
// The user is not allowed to send messages
if (error.status === 403) {
this.off('outgoingMessage', this.onPollOutgoingMessage.bind(this));
this.canEdit = false;
}
}
}
}
/**
* EventSource is a API for opening an HTTP
* connection for receiving push notifications
* from a server in real-time.
* We use it to sync the document with the server
*/
protected initCollaborationSSE() {
if (!this.isWebsocketFailed) {
return;
}
this.sse = new EventSource(this.toPollUrl('message'), {
withCredentials: true,
});
this.sse.onmessage = (event) => {
const { updatedDoc64, stateFingerprint, awareness64 } = JSON.parse(
// eslint-disable-next-line @typescript-eslint/no-unsafe-argument
event.data,
) as {
updatedDoc64?: string;
stateFingerprint?: string;
awareness64?: string;
};
if (awareness64) {
if (this.isDuplicateAwareness(awareness64)) {
return;
}
this.treatedAwarenessMessage = awareness64;
const awareness = Buffer.from(awareness64, 'base64');
this.onMessage({
data: awareness,
} as MessageEvent);
}
if (updatedDoc64) {
this.document.transact(() => {
Y.applyUpdate(this.document, Buffer.from(updatedDoc64, 'base64'));
}, this);
}
const localStateFingerprint = this.getStateFingerprint(this.document);
if (localStateFingerprint !== stateFingerprint) {
void this.pollSync();
} else {
this.seemsUnsyncCount = 0;
}
};
this.sse.onopen = () => {};
this.sse.onerror = (err) => {
console.error('SSE error:', err);
this.sse?.close();
setTimeout(() => {
this.initCollaborationSSE();
}, 5000);
};
}
/**
* Sync the document with the server.
*
* In some rare cases, the document may be out of sync.
* We use a fingerprint to compare documents,
* it happens that the local fingerprint is different from the server one
* when awareness plus the document are updated quickly.
* The system is resilient to this kind of problems, so `seemsUnsyncCount` should
* go back to 0 after a few seconds. If not, we will force a sync.
*/
public async pollSync(forseSync = false) {
if (!this.isWebsocketFailed || this.isSyncing) {
return;
}
this.seemsUnsyncCount++;
if (this.seemsUnsyncCount < this.seemsUnsyncMaxCount && !forseSync) {
return;
}
this.isSyncing = true;
try {
const { syncDoc64 } = await postPollSyncRequest({
pollUrl: this.toPollUrl('sync'),
localDoc64: toBase64(Y.encodeStateAsUpdate(this.document)),
});
if (syncDoc64) {
const uint8Array = Buffer.from(syncDoc64, 'base64');
Y.applyUpdate(this.document, uint8Array);
this.seemsUnsyncCount = 0;
}
} catch (error) {
console.error('Polling sync failed:', error);
} finally {
this.isSyncing = false;
}
}
/**
* Create a hash SHA-256 of the state vector of the document.
* Usefull to compare the state of the document.
* @param doc
* @returns
*/
public getStateFingerprint(doc: Y.Doc): string {
const stateVector = Y.encodeStateVector(doc);
return crypto.createHash('sha256').update(stateVector).digest('base64');
}
}
@@ -0,0 +1,179 @@
import { WebSocketStatus } from '@hocuspocus/provider';
import fetchMock from 'fetch-mock';
import * as Y from 'yjs';
if (typeof EventSource === 'undefined') {
const mockEventSource = jest.fn();
class MockEventSource {
constructor(...args: any[]) {
return mockEventSource(...args);
}
}
(global as any).EventSource = MockEventSource;
}
import { CollaborationProvider } from '../CollaborationProvider';
const mockApplyUpdate = jest.fn();
jest.mock('yjs', () => ({
...jest.requireActual('yjs'),
applyUpdate: (...args: any) => mockApplyUpdate(...args),
}));
describe('CollaborationProvider', () => {
let config: any;
let provider: CollaborationProvider;
let fakeWebsocketProvider: any;
beforeEach(() => {
fakeWebsocketProvider = {
on: jest.fn(),
open: jest.fn(),
attach: jest.fn(),
};
config = {
name: 'test',
url: 'ws://localhost/ws/',
canEdit: true,
websocketProvider: fakeWebsocketProvider,
};
provider = new CollaborationProvider(config);
});
afterEach(() => {
jest.clearAllMocks();
fetchMock.restore();
});
test('constructor initializes properties and attaches event handlers', () => {
expect(provider.canEdit).toBe(true);
expect((provider as any).url).toBe('ws://localhost/ws/');
expect(fakeWebsocketProvider.on).toHaveBeenCalled();
});
test('getStateFingerprint returns a consistent hash', () => {
const fingerprint1 = provider.getStateFingerprint(provider.document);
const fingerprint2 = provider.getStateFingerprint(provider.document);
expect(typeof fingerprint1).toBe('string');
expect(fingerprint1).toBe(fingerprint2);
});
test('onPollOutgoingMessage does nothing when websocket is not failed', async () => {
fetchMock.post(/http:\/\/localhost\/ws\/poll\/message\/.*/, {
body: JSON.stringify({ updated: false }),
});
provider.isWebsocketFailed = false;
const dummyMessage = {
toUint8Array: () => new Uint8Array([1, 2, 3]),
} as any;
await provider.onPollOutgoingMessage({ message: dummyMessage });
expect(fetchMock.called()).toBe(false);
});
test('onPollOutgoingMessage calls pollOutgoingMessageRequest and pollSync when updated is false', async () => {
provider.isWebsocketFailed = true;
const dummyMessage = {
toUint8Array: () => new Uint8Array([4, 5, 6]),
} as any;
fetchMock.post(/http:\/\/localhost\/ws\/poll\/message\/.*/, {
body: JSON.stringify({ updated: false }),
});
const pollSyncSpy = jest.spyOn(provider, 'pollSync').mockResolvedValue();
await provider.onPollOutgoingMessage({ message: dummyMessage });
expect(fetchMock.lastUrl()).toContain('http://localhost/ws/poll/message/');
expect(pollSyncSpy).toHaveBeenCalled();
});
test('onPollOutgoingMessage disables editing (canEdit becomes false) if API returns a 403 error', async () => {
provider.isWebsocketFailed = true;
const dummyMessage = {
toUint8Array: () => new Uint8Array([7, 8, 9]),
} as any;
fetchMock.post(/http:\/\/localhost\/ws\/poll\/message\/.*/, {
status: 403,
body: JSON.stringify({}),
});
// Stub the off method (inherited from event emitter) to observe its call.
provider.off = jest.fn();
await provider.onPollOutgoingMessage({ message: dummyMessage });
expect(fetchMock.lastUrl()).toContain('http://localhost/ws/poll/message/');
expect(provider.off).toHaveBeenCalled();
expect(provider.canEdit).toBe(false);
});
test('pollSync does nothing if websocket is not failed', async () => {
fetchMock.post(/http:\/\/localhost\/ws\/poll\/sync\/.*/, {
body: JSON.stringify({ syncDoc64: '123456' }),
});
provider.isWebsocketFailed = false;
await provider.pollSync();
expect(fetchMock.called()).toBe(false);
});
test('pollSync calls postPollSyncRequest when unsync count threshold is reached', async () => {
const update = Y.encodeStateAsUpdate(provider.document);
const syncDoc64 = Buffer.from(update).toString('base64');
fetchMock.post(/http:\/\/localhost\/ws\/poll\/sync\/.*/, {
body: JSON.stringify({ syncDoc64 }),
});
provider.isWebsocketFailed = true;
provider.seemsUnsyncCount = provider.seemsUnsyncMaxCount - 1;
await provider.pollSync();
const uint8Array = Buffer.from(syncDoc64, 'base64');
expect(mockApplyUpdate).toHaveBeenCalledWith(provider.document, uint8Array);
});
describe('onStatus', () => {
beforeEach(() => {
jest.useFakeTimers();
});
afterEach(() => {
jest.useRealTimers();
});
test('sets websocket failed and schedules polling on Connecting', () => {
const initPollingSpy = jest
.spyOn(provider, 'initPolling')
.mockImplementation(() => {});
const superOnStatusSpy = jest.spyOn(
Object.getPrototypeOf(provider),
'onStatus',
);
provider.onStatus({ status: WebSocketStatus.Connecting });
expect(provider.isWebsocketFailed).toBe(true);
// Fast-forward timer to trigger the scheduled initPolling call.
jest.runAllTimers();
expect(initPollingSpy).toHaveBeenCalled();
expect(superOnStatusSpy).toHaveBeenCalledWith({
status: WebSocketStatus.Connecting,
});
});
test('calls setPollDefaultValues on Connected', () => {
const setPollDefaultValuesSpy = jest
.spyOn(provider, 'setPollDefaultValues')
.mockImplementation(() => {});
const superOnStatusSpy = jest.spyOn(
Object.getPrototypeOf(provider),
'onStatus',
);
provider.onStatus({ status: WebSocketStatus.Connected });
expect(setPollDefaultValuesSpy).toHaveBeenCalled();
expect(superOnStatusSpy).toHaveBeenCalledWith({
status: WebSocketStatus.Connected,
});
});
});
});
@@ -4,10 +4,13 @@ import { create } from 'zustand';
import { Base64 } from '@/features/docs/doc-management';
import { CollaborationProvider } from '../libs/CollaborationProvider';
export interface UseCollaborationStore {
createProvider: (
providerUrl: string,
storeId: string,
canEdit: boolean,
initialDoc?: Base64,
) => HocuspocusProvider;
destroyProvider: () => void;
@@ -20,7 +23,7 @@ const defaultValues = {
export const useProviderStore = create<UseCollaborationStore>((set, get) => ({
...defaultValues,
createProvider: (wsUrl, storeId, initialDoc) => {
createProvider: (wsUrl, storeId, canEdit, initialDoc) => {
const doc = new Y.Doc({
guid: storeId,
});
@@ -29,10 +32,11 @@ export const useProviderStore = create<UseCollaborationStore>((set, get) => ({
Y.applyUpdate(doc, Buffer.from(initialDoc, 'base64'));
}
const provider = new HocuspocusProvider({
const provider = new CollaborationProvider({
url: wsUrl,
name: storeId,
document: doc,
canEdit,
});
set({
@@ -63,7 +63,7 @@ const DocPage = ({ id }: DocProps) => {
const { addTask } = useBroadcastStore();
const queryClient = useQueryClient();
const { replace } = useRouter();
useCollaboration(doc?.id, doc?.content);
useCollaboration(doc?.id, doc?.content, doc?.abilities.partial_update);
useEffect(() => {
if (doc?.title) {
+7
View File
@@ -5281,6 +5281,13 @@
dependencies:
"@types/node" "*"
"@types/[email protected]":
version "8.5.13"
resolved "https://registry.yarnpkg.com/@types/ws/-/ws-8.5.13.tgz#6414c280875e2691d0d1e080b05addbf5cb91e20"
integrity sha512-osM/gWBTPKgHV8XkTunnegTRIsvF6owmf5w+JtAfOw472dptdm0dlGv4xCt6GwQRcC2XVOvvRE/0bAoQcL2QkA==
dependencies:
"@types/node" "*"
"@types/yargs-parser@*":
version "21.0.3"
resolved "https://registry.yarnpkg.com/@types/yargs-parser/-/yargs-parser-21.0.3.tgz#815e30b786d2e8f0dcd85fd5bcf5e1a04d008f15"