chore: drop old client support (#14369)
This commit is contained in:
@@ -9,6 +9,7 @@ import {
|
||||
WebSocketServer,
|
||||
} from '@nestjs/websockets';
|
||||
import { ClsInterceptor } from 'nestjs-cls';
|
||||
import semver from 'semver';
|
||||
import { type Server, Socket } from 'socket.io';
|
||||
|
||||
import {
|
||||
@@ -49,10 +50,10 @@ type EventResponse<Data = any> = Data extends never
|
||||
data: Data;
|
||||
};
|
||||
|
||||
// sync-019: legacy 0.19.x clients (broadcast-doc-updates/push-doc-updates).
|
||||
// Remove after 2026-06-30 once metrics show 0 usage for 30 days.
|
||||
// 020+: receives space:broadcast-doc-updates (batch) and sends space:push-doc-update.
|
||||
type RoomType = 'sync' | `${string}:awareness` | 'sync-019';
|
||||
// sync: shared room for space membership checks and non-protocol broadcasts.
|
||||
// sync-025: legacy 0.25 doc sync protocol (space:broadcast-doc-update).
|
||||
// sync-026: current doc sync protocol (space:broadcast-doc-updates).
|
||||
type RoomType = 'sync' | 'sync-025' | 'sync-026' | `${string}:awareness`;
|
||||
|
||||
function Room(
|
||||
spaceId: string,
|
||||
@@ -61,6 +62,25 @@ function Room(
|
||||
return `${spaceId}:${type}`;
|
||||
}
|
||||
|
||||
const MIN_WS_CLIENT_VERSION = new semver.Range('>=0.25.0', {
|
||||
includePrerelease: true,
|
||||
});
|
||||
const DOC_UPDATES_PROTOCOL_026 = new semver.Range('>=0.26.0-0', {
|
||||
includePrerelease: true,
|
||||
});
|
||||
|
||||
type SyncProtocolRoomType = Extract<RoomType, 'sync-025' | 'sync-026'>;
|
||||
|
||||
function isSupportedWsClientVersion(clientVersion: string): boolean {
|
||||
return Boolean(
|
||||
semver.valid(clientVersion) && MIN_WS_CLIENT_VERSION.test(clientVersion)
|
||||
);
|
||||
}
|
||||
|
||||
function getSyncProtocolRoomType(clientVersion: string): SyncProtocolRoomType {
|
||||
return DOC_UPDATES_PROTOCOL_026.test(clientVersion) ? 'sync-026' : 'sync-025';
|
||||
}
|
||||
|
||||
enum SpaceType {
|
||||
Workspace = 'workspace',
|
||||
Userspace = 'userspace',
|
||||
@@ -90,16 +110,6 @@ interface LeaveSpaceAwarenessMessage {
|
||||
docId: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* @deprecated
|
||||
*/
|
||||
interface PushDocUpdatesMessage {
|
||||
spaceType: SpaceType;
|
||||
spaceId: string;
|
||||
docId: string;
|
||||
updates: string[];
|
||||
}
|
||||
|
||||
interface PushDocUpdateMessage {
|
||||
spaceType: SpaceType;
|
||||
spaceId: string;
|
||||
@@ -117,6 +127,15 @@ interface BroadcastDocUpdatesMessage {
|
||||
compressed?: boolean;
|
||||
}
|
||||
|
||||
interface BroadcastDocUpdateMessage {
|
||||
spaceType: SpaceType;
|
||||
spaceId: string;
|
||||
docId: string;
|
||||
update: string;
|
||||
timestamp: number;
|
||||
editor: string;
|
||||
}
|
||||
|
||||
interface LoadDocMessage {
|
||||
spaceType: SpaceType;
|
||||
spaceId: string;
|
||||
@@ -225,6 +244,11 @@ export class SpaceSyncGateway
|
||||
}
|
||||
}
|
||||
|
||||
private rejectJoin(client: Socket) {
|
||||
// Give socket.io a chance to flush the ack packet before disconnecting.
|
||||
setImmediate(() => client.disconnect());
|
||||
}
|
||||
|
||||
handleConnection() {
|
||||
this.connectionCount++;
|
||||
this.logger.debug(`New connection, total: ${this.connectionCount}`);
|
||||
@@ -252,23 +276,21 @@ export class SpaceSyncGateway
|
||||
return;
|
||||
}
|
||||
|
||||
const room025 = `${spaceType}:${Room(spaceId, 'sync-025')}`;
|
||||
const encodedUpdates = this.encodeUpdates(updates);
|
||||
|
||||
this.server
|
||||
.to(Room(spaceId, 'sync-019'))
|
||||
.emit('space:broadcast-doc-updates', {
|
||||
spaceType,
|
||||
for (const update of encodedUpdates) {
|
||||
const payload: BroadcastDocUpdateMessage = {
|
||||
spaceType: spaceType as SpaceType,
|
||||
spaceId,
|
||||
docId,
|
||||
updates: encodedUpdates,
|
||||
update,
|
||||
timestamp,
|
||||
editor,
|
||||
});
|
||||
metrics.socketio
|
||||
.counter('sync_019_broadcast')
|
||||
.add(encodedUpdates.length, { event: 'doc_updates_pushed' });
|
||||
editor: editor ?? '',
|
||||
};
|
||||
this.server.to(room025).emit('space:broadcast-doc-update', payload);
|
||||
}
|
||||
|
||||
const room = `${spaceType}:${Room(spaceId)}`;
|
||||
const room026 = `${spaceType}:${Room(spaceId, 'sync-026')}`;
|
||||
const payload = this.buildBroadcastPayload(
|
||||
spaceType as SpaceType,
|
||||
spaceId,
|
||||
@@ -277,7 +299,7 @@ export class SpaceSyncGateway
|
||||
timestamp,
|
||||
editor
|
||||
);
|
||||
this.server.to(room).emit('space:broadcast-doc-updates', payload);
|
||||
this.server.to(room026).emit('space:broadcast-doc-updates', payload);
|
||||
metrics.socketio
|
||||
.counter('doc_updates_broadcast')
|
||||
.add(payload.updates.length, {
|
||||
@@ -314,16 +336,34 @@ export class SpaceSyncGateway
|
||||
@MessageBody()
|
||||
{ spaceType, spaceId, clientVersion }: JoinSpaceMessage
|
||||
): Promise<EventResponse<{ clientId: string; success: boolean }>> {
|
||||
if (
|
||||
![SpaceType.Userspace, SpaceType.Workspace].includes(spaceType) ||
|
||||
/^0.1/.test(clientVersion)
|
||||
) {
|
||||
if (![SpaceType.Userspace, SpaceType.Workspace].includes(spaceType)) {
|
||||
this.rejectJoin(client);
|
||||
return { data: { clientId: client.id, success: false } };
|
||||
} else {
|
||||
if (spaceType === SpaceType.Workspace) {
|
||||
this.event.emit('workspace.embedding', { workspaceId: spaceId });
|
||||
}
|
||||
await this.selectAdapter(client, spaceType).join(user.id, spaceId);
|
||||
}
|
||||
|
||||
if (!isSupportedWsClientVersion(clientVersion)) {
|
||||
this.rejectJoin(client);
|
||||
return { data: { clientId: client.id, success: false } };
|
||||
}
|
||||
|
||||
if (spaceType === SpaceType.Workspace) {
|
||||
this.event.emit('workspace.embedding', { workspaceId: spaceId });
|
||||
}
|
||||
|
||||
const adapter = this.selectAdapter(client, spaceType);
|
||||
await adapter.join(user.id, spaceId);
|
||||
|
||||
const protocolRoomType = getSyncProtocolRoomType(clientVersion);
|
||||
const protocolRoom = adapter.room(spaceId, protocolRoomType);
|
||||
const otherProtocolRoom = adapter.room(
|
||||
spaceId,
|
||||
protocolRoomType === 'sync-025' ? 'sync-026' : 'sync-025'
|
||||
);
|
||||
if (client.rooms.has(otherProtocolRoom)) {
|
||||
await client.leave(otherProtocolRoom);
|
||||
}
|
||||
if (!client.rooms.has(protocolRoom)) {
|
||||
await client.join(protocolRoom);
|
||||
}
|
||||
|
||||
return { data: { clientId: client.id, success: true } };
|
||||
@@ -380,68 +420,8 @@ export class SpaceSyncGateway
|
||||
}
|
||||
|
||||
/**
|
||||
* @deprecated use [space:push-doc-update] instead, client should always merge updates on their own
|
||||
*
|
||||
* only 0.19.x client will send this event
|
||||
* client should always merge updates on their own
|
||||
*/
|
||||
@SubscribeMessage('space:push-doc-updates')
|
||||
async onReceiveDocUpdates(
|
||||
@ConnectedSocket() client: Socket,
|
||||
@CurrentUser() user: CurrentUser,
|
||||
@MessageBody()
|
||||
message: PushDocUpdatesMessage
|
||||
): Promise<EventResponse<{ accepted: true; timestamp?: number }>> {
|
||||
const { spaceType, spaceId, docId, updates } = message;
|
||||
const adapter = this.selectAdapter(client, spaceType);
|
||||
const id = new DocID(docId, spaceId);
|
||||
|
||||
// TODO(@forehalo): enable after frontend supporting doc revert
|
||||
// await this.ac.user(user.id).doc(spaceId, id.guid).assert('Doc.Update');
|
||||
const timestamp = await adapter.push(
|
||||
spaceId,
|
||||
id.guid,
|
||||
updates.map(update => Buffer.from(update, 'base64')),
|
||||
user.id
|
||||
);
|
||||
|
||||
metrics.socketio
|
||||
.counter('sync_019_event')
|
||||
.add(1, { event: 'push-doc-updates' });
|
||||
|
||||
// broadcast to 0.19.x clients
|
||||
client.to(Room(spaceId, 'sync-019')).emit('space:broadcast-doc-updates', {
|
||||
...message,
|
||||
timestamp,
|
||||
editor: user.id,
|
||||
});
|
||||
|
||||
// broadcast to new clients
|
||||
const decodedUpdates = updates.map(update => Buffer.from(update, 'base64'));
|
||||
const payload = this.buildBroadcastPayload(
|
||||
spaceType,
|
||||
spaceId,
|
||||
docId,
|
||||
decodedUpdates,
|
||||
timestamp,
|
||||
user.id
|
||||
);
|
||||
client
|
||||
.to(adapter.room(spaceId))
|
||||
.emit('space:broadcast-doc-updates', payload);
|
||||
metrics.socketio
|
||||
.counter('doc_updates_broadcast')
|
||||
.add(payload.updates.length, {
|
||||
mode: payload.compressed ? 'compressed' : 'batch',
|
||||
});
|
||||
|
||||
return {
|
||||
data: {
|
||||
accepted: true,
|
||||
timestamp,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@SubscribeMessage('space:push-doc-update')
|
||||
async onReceiveDocUpdate(
|
||||
@ConnectedSocket() client: Socket,
|
||||
@@ -461,16 +441,6 @@ export class SpaceSyncGateway
|
||||
user.id
|
||||
);
|
||||
|
||||
// broadcast to 0.19.x clients
|
||||
client.to(Room(spaceId, 'sync-019')).emit('space:broadcast-doc-updates', {
|
||||
spaceType,
|
||||
spaceId,
|
||||
docId,
|
||||
updates: [update],
|
||||
timestamp,
|
||||
editor: user.id,
|
||||
});
|
||||
|
||||
const payload = this.buildBroadcastPayload(
|
||||
spaceType,
|
||||
spaceId,
|
||||
@@ -480,7 +450,7 @@ export class SpaceSyncGateway
|
||||
user.id
|
||||
);
|
||||
client
|
||||
.to(adapter.room(spaceId))
|
||||
.to(adapter.room(spaceId, 'sync-026'))
|
||||
.emit('space:broadcast-doc-updates', payload);
|
||||
metrics.socketio
|
||||
.counter('doc_updates_broadcast')
|
||||
@@ -488,6 +458,17 @@ export class SpaceSyncGateway
|
||||
mode: payload.compressed ? 'compressed' : 'batch',
|
||||
});
|
||||
|
||||
client
|
||||
.to(adapter.room(spaceId, 'sync-025'))
|
||||
.emit('space:broadcast-doc-update', {
|
||||
spaceType,
|
||||
spaceId,
|
||||
docId,
|
||||
update,
|
||||
timestamp,
|
||||
editor: user.id,
|
||||
} satisfies BroadcastDocUpdateMessage);
|
||||
|
||||
return {
|
||||
data: {
|
||||
accepted: true,
|
||||
@@ -516,8 +497,18 @@ export class SpaceSyncGateway
|
||||
@ConnectedSocket() client: Socket,
|
||||
@CurrentUser() user: CurrentUser,
|
||||
@MessageBody()
|
||||
{ spaceType, spaceId, docId }: JoinSpaceAwarenessMessage
|
||||
{ spaceType, spaceId, docId, clientVersion }: JoinSpaceAwarenessMessage
|
||||
) {
|
||||
if (![SpaceType.Userspace, SpaceType.Workspace].includes(spaceType)) {
|
||||
this.rejectJoin(client);
|
||||
return { data: { clientId: client.id, success: false } };
|
||||
}
|
||||
|
||||
if (!isSupportedWsClientVersion(clientVersion)) {
|
||||
this.rejectJoin(client);
|
||||
return { data: { clientId: client.id, success: false } };
|
||||
}
|
||||
|
||||
await this.selectAdapter(client, spaceType).join(
|
||||
user.id,
|
||||
spaceId,
|
||||
@@ -555,13 +546,6 @@ export class SpaceSyncGateway
|
||||
.to(adapter.room(spaceId, roomType))
|
||||
.emit('space:collect-awareness', { spaceType, spaceId, docId });
|
||||
|
||||
// TODO(@forehalo): remove backward compatibility
|
||||
if (spaceType === SpaceType.Workspace) {
|
||||
client
|
||||
.to(adapter.room(spaceId, roomType))
|
||||
.emit('new-client-awareness-init');
|
||||
}
|
||||
|
||||
return { data: { clientId: client.id } };
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user