diff --git a/apps/meteor/ee/server/hooks/federation/index.ts b/apps/meteor/ee/server/hooks/federation/index.ts index b2c617c487113..cd9543b1a75c9 100644 --- a/apps/meteor/ee/server/hooks/federation/index.ts +++ b/apps/meteor/ee/server/hooks/federation/index.ts @@ -1,8 +1,7 @@ -import { api, FederationMatrix } from '@rocket.chat/core-services'; +import { FederationMatrix } from '@rocket.chat/core-services'; import { isEditedMessage, type IMessage, type IRoom, type IUser } from '@rocket.chat/core-typings'; import { MatrixBridgedRoom, Rooms } from '@rocket.chat/models'; -import notifications from '../../../../app/notifications/server/lib/Notifications'; import { callbacks } from '../../../../lib/callbacks'; import { afterLeaveRoomCallback } from '../../../../lib/callbacks/afterLeaveRoomCallback'; import { afterRemoveFromRoomCallback } from '../../../../lib/callbacks/afterRemoveFromRoomCallback'; @@ -17,8 +16,6 @@ import { FederationActions } from '../../../../server/services/room/hooks/Before callbacks.add('federation.afterCreateFederatedRoom', async (room, { owner, originalMemberList: members, options }) => { if (FederationActions.shouldPerformFederationAction(room)) { const federatedRoomId = options?.federatedRoomId; - // TODO: move this to the hooks folder - setupTypingEventListenerForRoom(room._id); if (!federatedRoomId) { // if room if exists, we don't want to create it again @@ -227,23 +224,3 @@ callbacks.add( callbacks.priority.HIGH, 'federation-matrix-after-create-direct-room', ); - -// TODO: THIS IS NOT READY FOR PRODUCTION! IMPOSSIBLE TO ADD ONE LISTENER PER ROOM! -const setupTypingEventListenerForRoom = (roomId: string): void => { - notifications.streamRoom.on(`${roomId}/user-activity`, (username, activity) => { - if (Array.isArray(activity) && (!activity.length || activity.includes('user-typing'))) { - void api.broadcast('user.typing', { - user: { username }, - isTyping: activity.includes('user-typing'), - roomId, - }); - } - }); -}; - -export const setupInternalEDUEventListeners = async () => { - const federatedRooms = await Rooms.findFederatedRooms({ projection: { _id: 1 } }).toArray(); - for (const room of federatedRooms) { - setupTypingEventListenerForRoom(room._id); - } -}; diff --git a/apps/meteor/ee/server/index.ts b/apps/meteor/ee/server/index.ts index 960b3b2f66245..e3604bbb36141 100644 --- a/apps/meteor/ee/server/index.ts +++ b/apps/meteor/ee/server/index.ts @@ -14,12 +14,6 @@ import './local-services/ldap/service'; import './methods/getReadReceipts'; import './patches'; import './hooks/federation'; -import { License } from '@rocket.chat/license'; export * from './apps/startup'; export { registerEEBroker } from './startup'; - -await License.onLicense('federation', async () => { - const { setupInternalEDUEventListeners } = await import('./hooks/federation'); - await setupInternalEDUEventListeners(); -}); diff --git a/apps/meteor/ee/server/startup/federation.ts b/apps/meteor/ee/server/startup/federation.ts index c98cf7b7c1371..12b4df1a8a370 100644 --- a/apps/meteor/ee/server/startup/federation.ts +++ b/apps/meteor/ee/server/startup/federation.ts @@ -5,6 +5,7 @@ import { License } from '@rocket.chat/license'; import { Logger } from '@rocket.chat/logger'; import { settings } from '../../../app/settings/server'; +import { StreamerCentral } from '../../../server/modules/streamer/streamer.module'; import { registerFederationRoutes } from '../api/federation'; const logger = new Logger('Federation'); @@ -27,6 +28,17 @@ export const startFederationService = async (): Promise => { logger.debug('Starting federation-matrix service'); federationMatrixService = await FederationMatrix.create(InstanceStatus.id()); + StreamerCentral.on('broadcast', (name, eventName, args) => { + if (!federationMatrixService) { + return; + } + if (name === 'notify-room' && eventName.endsWith('user-activity')) { + const [rid] = eventName.split('/'); + const [user, activity] = args; + void federationMatrixService.notifyUserTyping(rid, user, activity.includes('user-typing')); + } + }); + try { api.registerService(federationMatrixService); await registerFederationRoutes(federationMatrixService); diff --git a/apps/meteor/server/modules/listeners/listeners.module.ts b/apps/meteor/server/modules/listeners/listeners.module.ts index 61c4e05a31277..6a02318045c3d 100644 --- a/apps/meteor/server/modules/listeners/listeners.module.ts +++ b/apps/meteor/server/modules/listeners/listeners.module.ts @@ -151,7 +151,28 @@ export class ListenersModule { notifications.notifyRoom(rid, 'videoconf', callId); }); - service.onEvent('presence.status', ({ user }) => this.handlePresence({ user }, notifications)); + service.onEvent('presence.status', ({ user }) => { + const { _id, username, name, status, statusText, roles } = user; + if (!status || !username) { + return; + } + + notifications.notifyUserInThisInstance(_id, 'userData', { + type: 'updated', + id: _id, + diff: { + status, + ...(statusText && { statusText }), + }, + unset: {}, + }); + + notifications.notifyLoggedInThisInstance('user-status', [_id, username, STATUS_MAP[status], statusText, name, roles]); + + if (_id) { + notifications.sendPresence(_id, username, STATUS_MAP[status], statusText); + } + }); service.onEvent('user.updateCustomStatus', (userStatus) => { notifications.notifyLoggedInThisInstance('updateCustomUserStatus', { @@ -159,12 +180,10 @@ export class ListenersModule { }); }); - service.onEvent('federation-matrix.user.typing', ({ isTyping, roomId, username }) => { - notifications.notifyRoom(roomId, 'user-activity', username, isTyping ? ['user-typing'] : []); + service.onEvent('user.activity', ({ isTyping, roomId, user }) => { + notifications.notifyRoom(roomId, 'user-activity', user, isTyping ? ['user-typing'] : []); }); - service.onEvent('federation-matrix.user.presence.status', ({ user }) => this.handlePresence({ user }, notifications)); - service.onEvent('watch.messages', async ({ message }) => { if (!message.rid) { return; @@ -491,30 +510,4 @@ export class ListenersModule { notifications.streamRoomMessage.emit(roomId, acknowledgeMessage); }); } - - private handlePresence( - { user }: { user: Pick }, - notifications: NotificationsModule, - ): void { - const { _id, username, name, status, statusText, roles } = user; - if (!status || !username) { - return; - } - - notifications.notifyUserInThisInstance(_id, 'userData', { - type: 'updated', - id: _id, - diff: { - status, - ...(statusText && { statusText }), - }, - unset: {}, - }); - - notifications.notifyLoggedInThisInstance('user-status', [_id, username, STATUS_MAP[status], statusText, name, roles]); - - if (_id) { - notifications.sendPresence(_id, username, STATUS_MAP[status], statusText); - } - } } diff --git a/ee/packages/federation-matrix/src/FederationMatrix.ts b/ee/packages/federation-matrix/src/FederationMatrix.ts index 8ad7ffd9c94cf..60e54d5384d9f 100644 --- a/ee/packages/federation-matrix/src/FederationMatrix.ts +++ b/ee/packages/federation-matrix/src/FederationMatrix.ts @@ -114,24 +114,7 @@ export class FederationMatrix extends ServiceClass implements IFederationMatrixS instance.homeserverServices = getAllServices(); MatrixMediaService.setHomeserverServices(instance.homeserverServices); instance.buildMatrixHTTPRoutes(); - instance.onEvent('user.typing', async ({ isTyping, roomId, user: { username } }): Promise => { - if (!roomId || !username) { - return; - } - const externalRoomId = await MatrixBridgedRoom.getExternalRoomId(roomId); - if (!externalRoomId) { - return; - } - const localUser = await Users.findOneByUsername(username, { projection: { _id: 1 } }); - if (!localUser) { - return; - } - const externalUserId = await MatrixBridgedUser.getExternalUserIdByLocalUserId(localUser._id); - if (!externalUserId) { - return; - } - void instance.homeserverServices.edu.sendTypingNotification(externalRoomId, externalUserId, isTyping); - }); + instance.onEvent( 'presence.status', async ({ user }: { user: Pick }): Promise => { @@ -1002,4 +985,24 @@ export class FederationMatrix extends ServiceClass implements IFederationMatrixS } await this.homeserverServices.room.setPowerLevelForUser(matrixRoomId, senderMatrixUserId, matrixUserId, powerLevel); } + + async notifyUserTyping(rid: string, user: string, isTyping: boolean) { + if (!rid || !user) { + return; + } + const externalRoomId = await MatrixBridgedRoom.getExternalRoomId(rid); + if (!externalRoomId) { + return; + } + const localUser = await Users.findOneByUsername(user, { projection: { _id: 1 } }); + if (!localUser) { + return; + } + const externalUserId = await MatrixBridgedUser.getExternalUserIdByLocalUserId(localUser._id); + if (!externalUserId) { + return; + } + + void this.homeserverServices.edu.sendTypingNotification(externalRoomId, externalUserId, isTyping); + } } diff --git a/ee/packages/federation-matrix/src/events/edu.ts b/ee/packages/federation-matrix/src/events/edu.ts index 3fab701075a3c..47ab6ae31e7db 100644 --- a/ee/packages/federation-matrix/src/events/edu.ts +++ b/ee/packages/federation-matrix/src/events/edu.ts @@ -28,8 +28,8 @@ export const edus = async (emitter: Emitter) => { return; } - void api.broadcast('federation-matrix.user.typing', { - username: user.username, + void api.broadcast('user.activity', { + user: user.username, isTyping: data.typing, roomId: matrixRoom, }); @@ -71,8 +71,9 @@ export const edus = async (emitter: Emitter) => { ); const { _id, username, statusText, roles, name } = user; - void api.broadcast('federation-matrix.user.presence.status', { + void api.broadcast('presence.status', { user: { status, _id, username, statusText, roles, name }, + previousStatus: undefined, }); logger.debug(`Updated presence for user ${matrixUser.uid} to ${status} from Matrix federation`); } catch (error) { diff --git a/packages/core-services/src/events/Events.ts b/packages/core-services/src/events/Events.ts index 765cc04f304c1..c534ca1b368a6 100644 --- a/packages/core-services/src/events/Events.ts +++ b/packages/core-services/src/events/Events.ts @@ -149,12 +149,7 @@ export type EventSignatures = { scope?: string; }): void; 'user.updateCustomStatus'(userStatus: Omit): void; - 'user.typing'(data: { user: Partial; isTyping: boolean; roomId: string }): void; - 'federation-matrix.user.typing'(data: { username: string; isTyping: boolean; roomId: string }): void; - 'federation-matrix.user.presence.status'(data: { - user: Pick; - previousStatus?: UserStatus; - }): void; + 'user.activity'(data: { user: string; isTyping: boolean; roomId: string }): void; 'user.video-conference'(data: { userId: IUser['_id']; action: string; diff --git a/packages/core-services/src/types/IFederationMatrixService.ts b/packages/core-services/src/types/IFederationMatrixService.ts index a1ea231df3535..9be6df23fb8bb 100644 --- a/packages/core-services/src/types/IFederationMatrixService.ts +++ b/packages/core-services/src/types/IFederationMatrixService.ts @@ -35,4 +35,5 @@ export interface IFederationMatrixService { role: 'moderator' | 'owner' | 'leader' | 'user', ): Promise; inviteUsersToRoom(room: IRoomFederated, usersUserName: string[], inviter: Pick): Promise; + notifyUserTyping(rid: string, user: string, isTyping: boolean): Promise; }