diff --git a/apps/meteor/app/api/server/v1/chat.js b/apps/meteor/app/api/server/v1/chat.js index ab623563eee79..76f042c8cbc0d 100644 --- a/apps/meteor/app/api/server/v1/chat.js +++ b/apps/meteor/app/api/server/v1/chat.js @@ -366,31 +366,6 @@ API.v1.addRoute( }, ); -API.v1.addRoute( - 'chat.getMessageReadReceipts', - { authRequired: true }, - { - async get() { - const { messageId } = this.queryParams; - if (!messageId) { - return API.v1.failure({ - error: "The required 'messageId' param is missing.", - }); - } - - try { - return API.v1.success({ - receipts: await Meteor.call('getReadReceipts', { messageId }), - }); - } catch (error) { - return API.v1.failure({ - error: error.message, - }); - } - }, - }, -); - API.v1.addRoute( 'chat.reportMessage', { authRequired: true }, diff --git a/apps/meteor/app/lib/server/startup/settings.ts b/apps/meteor/app/lib/server/startup/settings.ts index b6c36e3230375..dc40a5b095e90 100644 --- a/apps/meteor/app/lib/server/startup/settings.ts +++ b/apps/meteor/app/lib/server/startup/settings.ts @@ -1179,6 +1179,23 @@ settingsRegistry.addGroup('Message', function () { public: true, }); }); + this.section('Read_Receipts', function () { + this.add('Message_Read_Receipt_Enabled', false, { + type: 'boolean', + enterprise: true, + invalidValue: false, + modules: ['message-read-receipt'], + public: true, + }); + this.add('Message_Read_Receipt_Store_Users', false, { + type: 'boolean', + enterprise: true, + invalidValue: false, + modules: ['message-read-receipt'], + public: true, + enableQuery: { _id: 'Message_Read_Receipt_Enabled', value: true }, + }); + }); this.add('Message_AllowEditing', true, { type: 'boolean', public: true, diff --git a/apps/meteor/app/models/server/models/Messages.js b/apps/meteor/app/models/server/models/Messages.js index 5b02eac244d85..0634f03bd9d51 100644 --- a/apps/meteor/app/models/server/models/Messages.js +++ b/apps/meteor/app/models/server/models/Messages.js @@ -1059,12 +1059,38 @@ export class Messages extends Base { return this.findOne(query, options); } - setAsRead(rid, until) { + setVisibleMessagesAsRead(rid, until) { return this.update( { rid, unread: true, ts: { $lt: until }, + $or: [ + { + tmid: { $exists: false }, + }, + { + tshow: true, + }, + ], + }, + { + $unset: { + unread: 1, + }, + }, + { + multi: true, + }, + ); + } + + setThreadMessagesAsRead(tmid, until) { + return this.update( + { + tmid, + unread: true, + ts: { $lt: until }, }, { $unset: { @@ -1090,10 +1116,37 @@ export class Messages extends Base { ); } - findUnreadMessagesByRoomAndDate(rid, after) { + findVisibleUnreadMessagesByRoomAndDate(rid, after) { const query = { unread: true, rid, + $or: [ + { + tmid: { $exists: false }, + }, + { + tshow: true, + }, + ], + }; + + if (after) { + query.ts = { $gt: after }; + } + + return this.find(query, { + fields: { + _id: 1, + }, + }); + } + + findUnreadThreadMessagesByDate(tmid, userId, after) { + const query = { + 'u._id': { $ne: userId }, + 'unread': true, + tmid, + 'tshow': { $exists: false }, }; if (after) { diff --git a/apps/meteor/app/threads/server/methods/getThreadMessages.js b/apps/meteor/app/threads/server/methods/getThreadMessages.js index 54d5f32b7691f..c92b631ad648f 100644 --- a/apps/meteor/app/threads/server/methods/getThreadMessages.js +++ b/apps/meteor/app/threads/server/methods/getThreadMessages.js @@ -1,5 +1,6 @@ import { Meteor } from 'meteor/meteor'; +import { callbacks } from '../../../../lib/callbacks'; import { Messages, Rooms } from '../../../models/server'; import { canAccessRoom } from '../../../authorization/server'; import { settings } from '../../../settings/server'; @@ -33,6 +34,11 @@ Meteor.methods({ throw new Meteor.Error('error-not-allowed', 'Not allowed', { method: 'getThreadMessages' }); } + if (!thread.tcount) { + return []; + } + + callbacks.run('beforeReadMessages', thread.rid, user._id); readThread({ userId: user._id, rid: thread.rid, tmid }); const result = Messages.findVisibleThreadByThreadId(tmid, { @@ -40,6 +46,7 @@ Meteor.methods({ ...(limit && { limit }), sort: { ts: -1 }, }).fetch(); + callbacks.runAsync('afterReadMessages', room._id, { uid: user._id, tmid }); return [thread, ...result]; }, diff --git a/apps/meteor/app/ui-cached-collection/client/models/CachedCollection.ts b/apps/meteor/app/ui-cached-collection/client/models/CachedCollection.ts index ba042a930bb87..ca2d3e26cc73d 100644 --- a/apps/meteor/app/ui-cached-collection/client/models/CachedCollection.ts +++ b/apps/meteor/app/ui-cached-collection/client/models/CachedCollection.ts @@ -44,7 +44,7 @@ export class CachedCollection extends Emitter<{ changed: T; re public eventType: EventType; - public version = 17; + public version = 18; public userRelated: boolean; diff --git a/apps/meteor/client/views/admin/info/LicenseCard.tsx b/apps/meteor/client/views/admin/info/LicenseCard.tsx index 177db8fb5835c..916fca4c2e09c 100644 --- a/apps/meteor/client/views/admin/info/LicenseCard.tsx +++ b/apps/meteor/client/views/admin/info/LicenseCard.tsx @@ -27,6 +27,7 @@ const LicenseCard = (): ReactElement => { const hasOmnichannel = modules.includes('livechat-enterprise'); const hasAuditing = modules.includes('auditing'); const hasCannedResponses = modules.includes('canned-responses'); + const hasReadReceipts = modules.includes('message-read-receipt'); const handleApplyLicense = useMutableCallback(() => setModal( @@ -64,6 +65,7 @@ const LicenseCard = (): ReactElement => { + )} diff --git a/apps/meteor/ee/app/license/server/bundles.ts b/apps/meteor/ee/app/license/server/bundles.ts index eedce00f4c86e..8eb80d4b8f019 100644 --- a/apps/meteor/ee/app/license/server/bundles.ts +++ b/apps/meteor/ee/app/license/server/bundles.ts @@ -13,7 +13,8 @@ export type BundleFeature = | 'device-management' | 'oauth-enterprise' | 'federation' - | 'videoconference-enterprise'; + | 'videoconference-enterprise' + | 'message-read-receipt'; interface IBundle { [key: string]: BundleFeature[]; @@ -36,6 +37,7 @@ const bundles: IBundle = { 'device-management', 'federation', 'videoconference-enterprise', + 'message-read-receipt', ], pro: [], }; diff --git a/apps/meteor/ee/app/message-read-receipt/server/hooks/afterReadMessages.ts b/apps/meteor/ee/app/message-read-receipt/server/hooks/afterReadMessages.ts new file mode 100644 index 0000000000000..bcbf0c751d389 --- /dev/null +++ b/apps/meteor/ee/app/message-read-receipt/server/hooks/afterReadMessages.ts @@ -0,0 +1,24 @@ +import type { IUser, IRoom, IMessage } from '@rocket.chat/core-typings'; +import { MessageReads } from '@rocket.chat/core-services'; + +import { ReadReceipt } from '../../../../server/lib/message-read-receipt/ReadReceipt'; +import { callbacks } from '../../../../../lib/callbacks'; +import { settings } from '../../../../../app/settings/server'; + +callbacks.add( + 'afterReadMessages', + (rid: IRoom['_id'], params: { uid: IUser['_id']; lastSeen?: Date; tmid?: IMessage['_id'] }) => { + if (!settings.get('Message_Read_Receipt_Enabled')) { + return; + } + const { uid, lastSeen, tmid } = params; + + if (tmid) { + MessageReads.readThread(uid, tmid); + } else if (lastSeen) { + ReadReceipt.markMessagesAsRead(rid, uid, lastSeen); + } + }, + callbacks.priority.MEDIUM, + 'message-read-receipt-afterReadMessages', +); diff --git a/apps/meteor/imports/message-read-receipt/server/hooks.js b/apps/meteor/ee/app/message-read-receipt/server/hooks/afterSaveMessage.ts similarity index 54% rename from apps/meteor/imports/message-read-receipt/server/hooks.js rename to apps/meteor/ee/app/message-read-receipt/server/hooks/afterSaveMessage.ts index 4155761a11890..278826dea4e70 100644 --- a/apps/meteor/imports/message-read-receipt/server/hooks.js +++ b/apps/meteor/ee/app/message-read-receipt/server/hooks/afterSaveMessage.ts @@ -1,17 +1,19 @@ import { Subscriptions } from '@rocket.chat/models'; +import type { IRoom, IMessage } from '@rocket.chat/core-typings'; +import { isEditedMessage, isOmnichannelRoom } from '@rocket.chat/core-typings'; -import { ReadReceipt } from './lib/ReadReceipt'; -import { callbacks } from '../../../lib/callbacks'; +import { ReadReceipt } from '../../../../server/lib/message-read-receipt/ReadReceipt'; +import { callbacks } from '../../../../../lib/callbacks'; callbacks.add( 'afterSaveMessage', - (message, room) => { + (message: IMessage, room: IRoom) => { // skips this callback if the message was edited - if (message.editedAt) { + if (isEditedMessage(message) && message.editedAt) { return message; } - if (room && !room.closedAt) { + if (!isOmnichannelRoom(room) || !room.closedAt) { // set subscription as read right after message was sent Promise.await(Subscriptions.setAsReadByRoomIdAndUserId(room._id, message.u._id)); } @@ -24,12 +26,3 @@ callbacks.add( callbacks.priority.MEDIUM, 'message-read-receipt-afterSaveMessage', ); - -callbacks.add( - 'afterReadMessages', - (rid, { uid, lastSeen }) => { - ReadReceipt.markMessagesAsRead(rid, uid, lastSeen); - }, - callbacks.priority.MEDIUM, - 'message-read-receipt-afterReadMessages', -); diff --git a/apps/meteor/ee/app/message-read-receipt/server/hooks/index.ts b/apps/meteor/ee/app/message-read-receipt/server/hooks/index.ts new file mode 100644 index 0000000000000..8ef3f27635d0a --- /dev/null +++ b/apps/meteor/ee/app/message-read-receipt/server/hooks/index.ts @@ -0,0 +1,2 @@ +import './afterReadMessages'; +import './afterSaveMessage'; diff --git a/apps/meteor/ee/app/message-read-receipt/server/index.ts b/apps/meteor/ee/app/message-read-receipt/server/index.ts new file mode 100644 index 0000000000000..89a8902672487 --- /dev/null +++ b/apps/meteor/ee/app/message-read-receipt/server/index.ts @@ -0,0 +1,5 @@ +import { onLicense } from '../../license/server'; + +onLicense('message-read-receipt', async () => { + await import('./hooks'); +}); diff --git a/apps/meteor/ee/client/startup/index.ts b/apps/meteor/ee/client/startup/index.ts index 496021c3f64a7..99382d37c5fe2 100644 --- a/apps/meteor/ee/client/startup/index.ts +++ b/apps/meteor/ee/client/startup/index.ts @@ -2,3 +2,4 @@ import './audit'; import './deviceManagement'; import './engagementDashboard'; import './slashCommands'; +import './readReceipt'; diff --git a/apps/meteor/ee/client/startup/readReceipt.ts b/apps/meteor/ee/client/startup/readReceipt.ts new file mode 100644 index 0000000000000..702d19ed57a3a --- /dev/null +++ b/apps/meteor/ee/client/startup/readReceipt.ts @@ -0,0 +1,34 @@ +import { Meteor } from 'meteor/meteor'; +import { Tracker } from 'meteor/tracker'; + +import { settings } from '../../../app/settings/client'; +import { MessageAction } from '../../../app/ui-utils/client'; +import { imperativeModal } from '../../../client/lib/imperativeModal'; +import { messageArgs } from '../../../client/lib/utils/messageArgs'; +import ReadReceiptsModal from '../../../client/views/room/modals/ReadReceiptsModal'; + +Meteor.startup(() => { + Tracker.autorun(() => { + const enabled = settings.get('Message_Read_Receipt_Store_Users'); + + if (!enabled) { + return MessageAction.removeButton('receipt-detail'); + } + + MessageAction.addButton({ + id: 'receipt-detail', + icon: 'info-circled', + label: 'Info', + context: ['starred', 'message', 'message-mobile', 'threads'], + action(_, props) { + const { message = messageArgs(this).msg } = props; + imperativeModal.open({ + component: ReadReceiptsModal, + props: { messageId: message._id, onClose: imperativeModal.close }, + }); + }, + order: 10, + group: 'menu', + }); + }); +}); diff --git a/apps/meteor/ee/definition/rest/index.ts b/apps/meteor/ee/definition/rest/index.ts index 059eecd93f89e..1f206f8e8c524 100644 --- a/apps/meteor/ee/definition/rest/index.ts +++ b/apps/meteor/ee/definition/rest/index.ts @@ -3,4 +3,5 @@ import './v1/engagementDashboard'; import './v1/omnichannel'; import './v1/sessions'; +import './v1/chat'; import './v1/roles'; diff --git a/apps/meteor/ee/definition/rest/v1/chat.ts b/apps/meteor/ee/definition/rest/v1/chat.ts new file mode 100644 index 0000000000000..dee7cadc65424 --- /dev/null +++ b/apps/meteor/ee/definition/rest/v1/chat.ts @@ -0,0 +1,34 @@ +import type { IMessage, ReadReceipt } from '@rocket.chat/core-typings'; +import Ajv from 'ajv'; + +const ajv = new Ajv({ + coerceTypes: true, +}); + +type GetMessageReadReceiptsProps = { + messageId: IMessage['_id']; +}; + +const getMessageReadReceiptsPropsSchema = { + type: 'object', + properties: { + messageId: { + type: 'string', + }, + }, + required: ['messageId'], + additionalProperties: false, +}; + +export const isGetMessageReadReceiptsProps = ajv.compile(getMessageReadReceiptsPropsSchema); + +declare module '@rocket.chat/rest-typings' { + // eslint-disable-next-line @typescript-eslint/naming-convention + interface Endpoints { + '/v1/chat.getMessageReadReceipts': { + GET: (params: GetMessageReadReceiptsProps) => { + receipts: ReadReceipt[]; + }; + }; + } +} diff --git a/apps/meteor/ee/server/api/chat.ts b/apps/meteor/ee/server/api/chat.ts new file mode 100644 index 0000000000000..238c0cb55d0b4 --- /dev/null +++ b/apps/meteor/ee/server/api/chat.ts @@ -0,0 +1,27 @@ +import { Meteor } from 'meteor/meteor'; + +import { API } from '../../../app/api/server/api'; +import { hasLicense } from '../../app/license/server/license'; + +API.v1.addRoute( + 'chat.getMessageReadReceipts', + { authRequired: true }, + { + async get() { + if (!hasLicense('message-read-receipt')) { + throw new Meteor.Error('error-action-not-allowed', 'This is an enterprise feature'); + } + + const { messageId } = this.queryParams; + if (!messageId) { + return API.v1.failure({ + error: "The required 'messageId' param is missing.", + }); + } + + return API.v1.success({ + receipts: await Meteor.call('getReadReceipts', { messageId }), + }); + }, + }, +); diff --git a/apps/meteor/ee/server/api/index.ts b/apps/meteor/ee/server/api/index.ts index f50324e3c6d3e..96dc64c5ced41 100644 --- a/apps/meteor/ee/server/api/index.ts +++ b/apps/meteor/ee/server/api/index.ts @@ -2,4 +2,5 @@ import './api'; import './ldap'; import './licenses'; import './sessions'; +import './chat'; import './roles'; diff --git a/apps/meteor/ee/server/index.ts b/apps/meteor/ee/server/index.ts index d7ad69be8830a..279a4bf028372 100644 --- a/apps/meteor/ee/server/index.ts +++ b/apps/meteor/ee/server/index.ts @@ -6,6 +6,7 @@ import '../app/api-enterprise/server/index'; import '../app/authorization/server/index'; import '../app/canned-responses/server/index'; import '../app/livechat-enterprise/server/index'; +import '../app/message-read-receipt/server/index'; import '../app/voip-enterprise/server/index'; import '../app/settings/server/index'; import '../app/teams-mention/server/index'; @@ -15,3 +16,4 @@ import './requestSeatsRoute'; import './configuration/index'; import './local-services/ldap/service'; import './settings/index'; +import './methods/getReadReceipts'; diff --git a/apps/meteor/ee/server/lib/constants.ts b/apps/meteor/ee/server/lib/constants.ts new file mode 100644 index 0000000000000..dfaa54af84d2a --- /dev/null +++ b/apps/meteor/ee/server/lib/constants.ts @@ -0,0 +1 @@ +export const MAX_ROOM_SIZE_CHECK_INDIVIDUAL_READ_RECEIPTS = 50; diff --git a/apps/meteor/imports/message-read-receipt/server/lib/ReadReceipt.js b/apps/meteor/ee/server/lib/message-read-receipt/ReadReceipt.js similarity index 81% rename from apps/meteor/imports/message-read-receipt/server/lib/ReadReceipt.js rename to apps/meteor/ee/server/lib/message-read-receipt/ReadReceipt.js index 910133e72ff0b..b0443b0ab2218 100644 --- a/apps/meteor/imports/message-read-receipt/server/lib/ReadReceipt.js +++ b/apps/meteor/ee/server/lib/message-read-receipt/ReadReceipt.js @@ -26,7 +26,7 @@ const updateMessages = debounceByRoomId( return; } - Messages.setAsRead(_id, firstSubscription.ls); + Messages.setVisibleMessagesAsRead(_id, firstSubscription.ls); if (lm <= firstSubscription.ls) { Rooms.setLastMessageAsRead(_id); @@ -47,7 +47,7 @@ export const ReadReceipt = { return; } - this.storeReadReceipts(Messages.findUnreadMessagesByRoomAndDate(roomId, userLastSeen), roomId, userId); + this.storeReadReceipts(Messages.findVisibleUnreadMessagesByRoomAndDate(roomId, userLastSeen), roomId, userId); updateMessages(room); }, @@ -71,6 +71,21 @@ export const ReadReceipt = { this.storeReadReceipts([{ _id: message._id }], roomId, userId, extraData); }, + storeThreadMessagesReadReceipts(tmid, userId, userLastSeen) { + if (!settings.get('Message_Read_Receipt_Enabled')) { + return; + } + + const message = Messages.findOneById(tmid, { fields: { tlm: 1, rid: 1 } }); + + // if users last seen is greater than thread's last message, it means the user has already marked this thread as read + if (!message || userLastSeen > message.tlm) { + return; + } + + this.storeReadReceipts(Messages.findUnreadThreadMessagesByDate(tmid, userId, userLastSeen), message.rid, userId); + }, + async storeReadReceipts(messages, roomId, userId, extraData = {}) { if (settings.get('Message_Read_Receipt_Store_Users')) { const ts = new Date(); diff --git a/apps/meteor/ee/server/local-services/message-reads/service.ts b/apps/meteor/ee/server/local-services/message-reads/service.ts new file mode 100644 index 0000000000000..8ed82c66af0c8 --- /dev/null +++ b/apps/meteor/ee/server/local-services/message-reads/service.ts @@ -0,0 +1,50 @@ +import { MessageReads, Subscriptions } from '@rocket.chat/models'; +import type { ISubscription } from '@rocket.chat/core-typings'; +import { ServiceClassInternal } from '@rocket.chat/core-services'; + +import type { IMessageReadsService } from '../../sdk/types/IMessageReadsService'; +import { Messages } from '../../../../app/models/server'; +import { ReadReceipt } from '../../lib/message-read-receipt/ReadReceipt'; +import { MAX_ROOM_SIZE_CHECK_INDIVIDUAL_READ_RECEIPTS } from '../../lib/constants'; + +export class MessageReadsService extends ServiceClassInternal implements IMessageReadsService { + protected name = 'message-reads'; + + async readThread(userId: string, tmid: string): Promise { + const read = await MessageReads.findOneByUserIdAndThreadId(userId, tmid); + + const threadMessage = Messages.findOneById(tmid, { projection: { ts: 1, tlm: 1, rid: 1 } }); + if (!threadMessage || !threadMessage.tlm) { + return; + } + + await MessageReads.updateReadTimestampByUserIdAndThreadId(userId, tmid); + ReadReceipt.storeThreadMessagesReadReceipts(tmid, userId, read?.ls || threadMessage.ts); + + // doesn't mark as read if not all room members have read the thread + const membersCount = await Subscriptions.countUnarchivedByRoomId(threadMessage.rid); + + if (membersCount <= MAX_ROOM_SIZE_CHECK_INDIVIDUAL_READ_RECEIPTS) { + const subscriptions = await Subscriptions.findUnarchivedByRoomId(threadMessage.rid, { + projection: { 'u._id': 1 }, + }).toArray(); + const members = subscriptions.map((s: ISubscription) => s.u._id); + + const totalMessageReads = await MessageReads.countByThreadAndUserIds(tmid, members); + if (totalMessageReads < membersCount) { + return; + } + } else { + // for large rooms, mark as read if there are as many reads as room members to improve performance (instead of checking each read) + const totalMessageReads = await MessageReads.countByThreadId(tmid); + if (totalMessageReads < membersCount) { + return; + } + } + + const firstRead = await MessageReads.getMinimumLastSeenByThreadId(tmid); + if (firstRead?.ls) { + Messages.setThreadMessagesAsRead(tmid, firstRead.ls); + } + } +} diff --git a/apps/meteor/imports/message-read-receipt/server/api/methods/getReadReceipts.js b/apps/meteor/ee/server/methods/getReadReceipts.js similarity index 63% rename from apps/meteor/imports/message-read-receipt/server/api/methods/getReadReceipts.js rename to apps/meteor/ee/server/methods/getReadReceipts.js index a037fea749108..9e20297d08e0d 100644 --- a/apps/meteor/imports/message-read-receipt/server/api/methods/getReadReceipts.js +++ b/apps/meteor/ee/server/methods/getReadReceipts.js @@ -1,12 +1,17 @@ import { Meteor } from 'meteor/meteor'; import { check } from 'meteor/check'; -import { Messages } from '../../../../../app/models/server'; -import { canAccessRoomId } from '../../../../../app/authorization/server'; -import { ReadReceipt } from '../../lib/ReadReceipt'; +import { Messages } from '../../../app/models/server'; +import { canAccessRoomId } from '../../../app/authorization/server'; +import { hasLicense } from '../../app/license/server/license'; +import { ReadReceipt } from '../lib/message-read-receipt/ReadReceipt'; Meteor.methods({ - async getReadReceipts({ messageId }) { + getReadReceipts({ messageId }) { + if (!hasLicense('message-read-receipt')) { + throw new Meteor.Error('error-action-not-allowed', 'This is an enterprise feature', { method: 'getReadReceipts' }); + } + if (!messageId) { throw new Meteor.Error('error-invalid-message', "The required 'messageId' param is missing.", { method: 'getReadReceipts' }); } diff --git a/apps/meteor/server/models/ReadReceipts.ts b/apps/meteor/ee/server/models/ReadReceipts.ts similarity index 76% rename from apps/meteor/server/models/ReadReceipts.ts rename to apps/meteor/ee/server/models/ReadReceipts.ts index d54a985ef18a2..079c7c34a6a80 100644 --- a/apps/meteor/server/models/ReadReceipts.ts +++ b/apps/meteor/ee/server/models/ReadReceipts.ts @@ -1,6 +1,6 @@ import { registerModel } from '@rocket.chat/models'; -import { db } from '../database/utils'; +import { db } from '../../../server/database/utils'; import { ReadReceiptsRaw } from './raw/ReadReceipts'; registerModel('IReadReceiptsModel', new ReadReceiptsRaw(db)); diff --git a/apps/meteor/server/models/raw/ReadReceipts.ts b/apps/meteor/ee/server/models/raw/ReadReceipts.ts similarity index 91% rename from apps/meteor/server/models/raw/ReadReceipts.ts rename to apps/meteor/ee/server/models/raw/ReadReceipts.ts index 9763a2616c138..bd5eb76bc112e 100644 --- a/apps/meteor/server/models/raw/ReadReceipts.ts +++ b/apps/meteor/ee/server/models/raw/ReadReceipts.ts @@ -2,7 +2,7 @@ import type { ReadReceipt, RocketChatRecordDeleted } from '@rocket.chat/core-typ import type { IReadReceiptsModel } from '@rocket.chat/model-typings'; import type { Collection, FindCursor, Db, IndexDescription } from 'mongodb'; -import { BaseRaw } from './BaseRaw'; +import { BaseRaw } from '../../../../server/models/raw/BaseRaw'; export class ReadReceiptsRaw extends BaseRaw implements IReadReceiptsModel { constructor(db: Db, trash?: Collection>) { diff --git a/apps/meteor/ee/server/models/startup.ts b/apps/meteor/ee/server/models/startup.ts index 7b803b4c95217..8a9b09d98e165 100644 --- a/apps/meteor/ee/server/models/startup.ts +++ b/apps/meteor/ee/server/models/startup.ts @@ -7,5 +7,6 @@ onLicense('livechat-enterprise', () => { import('./LivechatUnit'); import('./LivechatUnitMonitors'); import('./LivechatRooms'); + import('./ReadReceipts'); import('./LivechatDepartment'); }); diff --git a/apps/meteor/ee/server/sdk/index.ts b/apps/meteor/ee/server/sdk/index.ts index 795af715544a9..ea6a92b2a9d87 100644 --- a/apps/meteor/ee/server/sdk/index.ts +++ b/apps/meteor/ee/server/sdk/index.ts @@ -2,6 +2,8 @@ import { proxifyWithWait } from '@rocket.chat/core-services'; import type { ILDAPEEService } from './types/ILDAPEEService'; import type { IInstanceService } from './types/IInstanceService'; +import type { IMessageReadsService } from './types/IMessageReadsService'; export const LDAPEE = proxifyWithWait('ldap-enterprise'); export const Instance = proxifyWithWait('instance'); +export const MessageReads = proxifyWithWait('message-reads'); diff --git a/apps/meteor/ee/server/sdk/types/IMessageReadsService.ts b/apps/meteor/ee/server/sdk/types/IMessageReadsService.ts new file mode 100644 index 0000000000000..ec7987bb8f6d3 --- /dev/null +++ b/apps/meteor/ee/server/sdk/types/IMessageReadsService.ts @@ -0,0 +1,3 @@ +export interface IMessageReadsService { + readThread(userId: string, threadId: string): Promise; +} diff --git a/apps/meteor/ee/server/sdk/types/index.ts b/apps/meteor/ee/server/sdk/types/index.ts new file mode 100644 index 0000000000000..ad33c3085a597 --- /dev/null +++ b/apps/meteor/ee/server/sdk/types/index.ts @@ -0,0 +1,2 @@ +export * from './ILDAPEEService'; +export * from './IMessageReadsService'; diff --git a/apps/meteor/ee/server/startup/services.ts b/apps/meteor/ee/server/startup/services.ts index bfaca103055b0..7ee7a7abc2be0 100644 --- a/apps/meteor/ee/server/startup/services.ts +++ b/apps/meteor/ee/server/startup/services.ts @@ -2,6 +2,7 @@ import { api } from '@rocket.chat/core-services'; import { EnterpriseSettings } from '../../app/settings/server/settings.internalService'; import { LDAPEEService } from '../local-services/ldap/service'; +import { MessageReadsService } from '../local-services/message-reads/service'; import { InstanceService } from '../local-services/instance/service'; import { LicenseService } from '../../app/license/server/license.internalService'; import { isRunningMs } from '../../../server/lib/isRunningMs'; @@ -10,6 +11,7 @@ import { isRunningMs } from '../../../server/lib/isRunningMs'; api.registerService(new EnterpriseSettings()); api.registerService(new LDAPEEService()); api.registerService(new LicenseService()); +api.registerService(new MessageReadsService()); // when not running micro services we want to start up the instance intercom if (!isRunningMs()) { diff --git a/apps/meteor/imports/message-read-receipt/server/index.js b/apps/meteor/imports/message-read-receipt/server/index.js deleted file mode 100644 index b58297d84fc66..0000000000000 --- a/apps/meteor/imports/message-read-receipt/server/index.js +++ /dev/null @@ -1,4 +0,0 @@ -import './hooks'; -import './settings'; - -import './api/methods/getReadReceipts'; diff --git a/apps/meteor/imports/message-read-receipt/server/settings.ts b/apps/meteor/imports/message-read-receipt/server/settings.ts deleted file mode 100644 index 8a20742847af3..0000000000000 --- a/apps/meteor/imports/message-read-receipt/server/settings.ts +++ /dev/null @@ -1,14 +0,0 @@ -import { settingsRegistry } from '../../../app/settings/server'; - -settingsRegistry.add('Message_Read_Receipt_Enabled', false, { - group: 'Message', - type: 'boolean', - public: true, -}); - -settingsRegistry.add('Message_Read_Receipt_Store_Users', false, { - group: 'Message', - type: 'boolean', - public: true, - enableQuery: { _id: 'Message_Read_Receipt_Enabled', value: true }, -}); diff --git a/apps/meteor/imports/startup/server/index.ts b/apps/meteor/imports/startup/server/index.ts index 8b7c5bd8020b1..90c02562d36ca 100644 --- a/apps/meteor/imports/startup/server/index.ts +++ b/apps/meteor/imports/startup/server/index.ts @@ -1,2 +1 @@ -import '../../message-read-receipt/server'; import '../../personal-access-tokens/server'; diff --git a/apps/meteor/lib/callbacks.ts b/apps/meteor/lib/callbacks.ts index 749d6d17193a6..312802fe7ce2e 100644 --- a/apps/meteor/lib/callbacks.ts +++ b/apps/meteor/lib/callbacks.ts @@ -44,7 +44,7 @@ type EventLikeCallbackSignatures = { 'afterDeleteMessage': (message: IMessage, room: IRoom) => void; 'validateUserRoles': (userData: Partial) => void; 'workspaceLicenseChanged': (license: string) => void; - 'afterReadMessages': (rid: IRoom['_id'], params: { uid: IUser['_id']; lastSeen: Date }) => void; + 'afterReadMessages': (rid: IRoom['_id'], params: { uid: IUser['_id']; lastSeen?: Date; tmid?: IMessage['_id'] }) => void; 'beforeReadMessages': (rid: IRoom['_id'], uid: IUser['_id']) => void; 'afterDeleteUser': (user: IUser) => void; 'afterFileUpload': (params: { user: IUser; room: IRoom; message: IMessage }) => void; diff --git a/apps/meteor/packages/rocketchat-i18n/i18n/en.i18n.json b/apps/meteor/packages/rocketchat-i18n/i18n/en.i18n.json index ed5de5db13150..13ea0c59b2bf9 100644 --- a/apps/meteor/packages/rocketchat-i18n/i18n/en.i18n.json +++ b/apps/meteor/packages/rocketchat-i18n/i18n/en.i18n.json @@ -3832,6 +3832,7 @@ "Reactions": "Reactions", "Read_by": "Read by", "Read_only": "Read Only", + "Read_Receipts": "Read Receipts", "This_room_is_read_only": "This room is read only", "Only_people_with_permission_can_send_messages_here": "Only people with permission can send messages here", "Read_only_changed_successfully": "Read only changed successfully", diff --git a/apps/meteor/server/methods/readThreads.js b/apps/meteor/server/methods/readThreads.js index 4a735ba9a47e2..afab0fe355e2f 100644 --- a/apps/meteor/server/methods/readThreads.js +++ b/apps/meteor/server/methods/readThreads.js @@ -5,6 +5,7 @@ import { settings } from '../../app/settings/server'; import { Messages, Rooms } from '../../app/models/server'; import { canAccessRoom } from '../../app/authorization/server'; import { readThread } from '../../app/threads/server/functions'; +import { callbacks } from '../../lib/callbacks'; Meteor.methods({ readThreads(tmid) { @@ -22,12 +23,15 @@ Meteor.methods({ } const user = Meteor.user(); + const room = Rooms.findOneById(thread.rid); if (!canAccessRoom(room, user)) { throw new Meteor.Error('error-not-allowed', 'Not allowed', { method: 'getThreadMessages' }); } - return readThread({ userId: user._id, rid: thread.rid, tmid }); + callbacks.run('beforeReadMessages', thread.rid, user._id); + readThread({ userId: user._id, rid: thread.rid, tmid }); + callbacks.runAsync('afterReadMessages', room._id, { uid: user._id, tmid }); }, }); diff --git a/apps/meteor/server/models/MessageReads.ts b/apps/meteor/server/models/MessageReads.ts new file mode 100644 index 0000000000000..43984e9e2b4a1 --- /dev/null +++ b/apps/meteor/server/models/MessageReads.ts @@ -0,0 +1,6 @@ +import { registerModel } from '@rocket.chat/models'; + +import { db } from '../database/utils'; +import { MessageReadsRaw } from './raw/MessageReads'; + +registerModel('IMessageReadsModel', new MessageReadsRaw(db)); diff --git a/apps/meteor/server/models/raw/MessageReads.ts b/apps/meteor/server/models/raw/MessageReads.ts new file mode 100644 index 0000000000000..bebe0abc2b703 --- /dev/null +++ b/apps/meteor/server/models/raw/MessageReads.ts @@ -0,0 +1,62 @@ +import type { MessageReads, IUser, IMessage, RocketChatRecordDeleted } from '@rocket.chat/core-typings'; +import type { IMessageReadsModel } from '@rocket.chat/model-typings'; +import type { Collection, Db, IndexDescription, UpdateResult } from 'mongodb'; + +import { BaseRaw } from './BaseRaw'; + +export class MessageReadsRaw extends BaseRaw implements IMessageReadsModel { + constructor(db: Db, trash?: Collection>) { + super(db, 'message_reads', trash); + } + + protected modelIndexes(): IndexDescription[] { + return [{ key: { tmid: 1, userId: 1 }, unique: true }, { key: { ls: 1 } }]; + } + + async findOneByUserIdAndThreadId(userId: IUser['_id'], tmid: IMessage['_id']): Promise { + return this.findOne({ userId, tmid }); + } + + getMinimumLastSeenByThreadId(tmid: IMessage['_id']): Promise { + return this.findOne( + { + tmid, + }, + { + sort: { + ls: 1, + }, + }, + ); + } + + updateReadTimestampByUserIdAndThreadId(userId: IUser['_id'], tmid: IMessage['_id']): Promise { + const query = { + userId, + tmid, + }; + + const update = { + $set: { + ls: new Date(), + }, + }; + + return this.updateOne(query, update, { upsert: true }); + } + + async countByThreadAndUserIds(tmid: IMessage['_id'], userIds: IUser['_id'][]): Promise { + const query = { + tmid, + userId: { $in: userIds }, + }; + return this.col.countDocuments(query); + } + + async countByThreadId(tmid: IMessage['_id']): Promise { + const query = { + tmid, + }; + return this.col.countDocuments(query); + } +} diff --git a/apps/meteor/server/models/raw/Subscriptions.ts b/apps/meteor/server/models/raw/Subscriptions.ts index d579c52eb4ec8..33e868bcd4596 100644 --- a/apps/meteor/server/models/raw/Subscriptions.ts +++ b/apps/meteor/server/models/raw/Subscriptions.ts @@ -71,6 +71,16 @@ export class SubscriptionsRaw extends BaseRaw implements ISubscri return this.find(query, options); } + findUnarchivedByRoomId(roomId: string, options: FindOptions = {}): FindCursor { + const query = { + 'rid': roomId, + 'archived': { $ne: true }, + 'u._id': { $exists: true }, + }; + + return this.find(query, options); + } + findByRoomIdAndNotUserId(roomId: string, userId: string, options: FindOptions = {}): FindCursor { const query = { 'rid': roomId, @@ -102,6 +112,15 @@ export class SubscriptionsRaw extends BaseRaw implements ISubscri return this.col.countDocuments(query); } + countUnarchivedByRoomId(rid: string): Promise { + const query = { + rid, + 'archived': { $ne: true }, + 'u._id': { $exists: true }, + }; + return this.col.countDocuments(query); + } + async isUserInRole(uid: IUser['_id'], roleId: IRole['_id'], rid?: IRoom['_id']): Promise { if (rid == null) { return null; diff --git a/apps/meteor/server/models/startup.ts b/apps/meteor/server/models/startup.ts index 32af7c43cd03b..d72c43cfe4a50 100644 --- a/apps/meteor/server/models/startup.ts +++ b/apps/meteor/server/models/startup.ts @@ -37,7 +37,7 @@ import './OEmbedCache'; import './PbxEvents'; import './PushToken'; import './Permissions'; -import './ReadReceipts'; +import './MessageReads'; import './Reports'; import './Roles'; import './Rooms'; diff --git a/apps/meteor/tests/end-to-end/api/05-chat.js b/apps/meteor/tests/end-to-end/api/05-chat.js index c46f4913f2a62..2c4c7a98366dd 100644 --- a/apps/meteor/tests/end-to-end/api/05-chat.js +++ b/apps/meteor/tests/end-to-end/api/05-chat.js @@ -1347,11 +1347,19 @@ describe('[Chat]', function () { }); describe('[/chat.getMessageReadReceipts]', () => { + const isEnterprise = typeof process.env.IS_EE === 'string' ? process.env.IS_EE === 'true' : !!process.env.IS_EE; describe('when execute successfully', () => { - it("should return the statusCode 200 and 'receipts' property and should be equal an array", (done) => { + it('should return statusCode: 200 and an array of receipts when running EE', function (done) { + if (!isEnterprise) { + this.skip(); + } + request - .get(api(`chat.getMessageReadReceipts?messageId=${message._id}`)) + .get(api(`chat.getMessageReadReceipts`)) .set(credentials) + .query({ + messageId: message._id, + }) .expect('Content-Type', 'application/json') .expect(200) .expect((res) => { @@ -1363,7 +1371,33 @@ describe('[Chat]', function () { }); describe('when an error occurs', () => { - it('should return statusCode 400 and an error', (done) => { + it('should throw error-action-not-allowed error when not running EE', function (done) { + // TODO this is not the right way to do it. We're doing this way for now just because we have separate CI jobs for EE and CE, + // ideally we should have a single CI job that adds a license and runs both CE and EE tests. + if (isEnterprise) { + this.skip(); + } + request + .get(api(`chat.getMessageReadReceipts`)) + .set(credentials) + .query({ + messageId: message._id, + }) + .expect('Content-Type', 'application/json') + .expect(400) + .expect((res) => { + expect(res.body).to.have.property('success', false); + expect(res.body).to.have.property('error', 'This is an enterprise feature [error-action-not-allowed]'); + expect(res.body).to.have.property('errorType', 'error-action-not-allowed'); + }) + .end(done); + }); + + it('should return statusCode: 400 and an error when no messageId is provided', function (done) { + if (!isEnterprise) { + this.skip(); + } + request .get(api('chat.getMessageReadReceipts')) .set(credentials) diff --git a/packages/core-services/src/index.ts b/packages/core-services/src/index.ts index e6b4c82aa271b..23f912ac056aa 100644 --- a/packages/core-services/src/index.ts +++ b/packages/core-services/src/index.ts @@ -23,6 +23,7 @@ import type { ITeamAutocompleteResult, IListRoomsFilter, } from './types/ITeamService'; +import type { IMessageReadsService } from './types/IMessageReadsService'; import type { IRoomService, ICreateRoomParams, ISubscriptionExtraData } from './types/IRoomService'; import type { IMediaService, ResizeResult } from './types/IMediaService'; import type { IVoipService } from './types/IVoipService'; @@ -76,6 +77,7 @@ export { IOmnichannelVoipService, IPresence, IPushService, + IMessageReadsService, IRoomService, ISAUMonitorService, ISubscriptionExtraData, @@ -121,6 +123,7 @@ export const Banner = proxifyWithWait('banner'); export const UiKitCoreApp = proxifyWithWait('uikit-core-app'); export const NPS = proxifyWithWait('nps'); export const Team = proxifyWithWait('team'); +export const MessageReads = proxifyWithWait('message-reads'); export const Room = proxifyWithWait('room'); export const Media = proxifyWithWait('media'); export const Voip = proxifyWithWait('voip'); diff --git a/packages/core-services/src/types/IMessageReadsService.ts b/packages/core-services/src/types/IMessageReadsService.ts new file mode 100644 index 0000000000000..6bb892abd3632 --- /dev/null +++ b/packages/core-services/src/types/IMessageReadsService.ts @@ -0,0 +1,5 @@ +import type { IServiceClass } from './ServiceClass'; + +export interface IMessageReadsService extends IServiceClass { + readThread(userId: string, threadId: string): Promise; +} diff --git a/packages/core-typings/src/MessageReads.ts b/packages/core-typings/src/MessageReads.ts new file mode 100644 index 0000000000000..c95cec5c5d170 --- /dev/null +++ b/packages/core-typings/src/MessageReads.ts @@ -0,0 +1,9 @@ +import type { IMessage } from './IMessage/IMessage'; +import type { IUser } from './IUser'; + +export type MessageReads = { + _id: string; + tmid: IMessage['_id']; + ls: Date; + userId: IUser['_id']; +}; diff --git a/packages/core-typings/src/index.ts b/packages/core-typings/src/index.ts index 4c509fcdbf0b1..56316aa3871b5 100644 --- a/packages/core-typings/src/index.ts +++ b/packages/core-typings/src/index.ts @@ -59,6 +59,7 @@ export * from './ICustomUserStatus'; export * from './IEmailMessageHistory'; export * from './ReadReceipt'; +export * from './MessageReads'; export * from './IUpload'; export * from './IOEmbedCache'; export * from './IOembed'; diff --git a/packages/model-typings/src/index.ts b/packages/model-typings/src/index.ts index 222df26698c31..a99cfd2c12ff3 100644 --- a/packages/model-typings/src/index.ts +++ b/packages/model-typings/src/index.ts @@ -42,6 +42,7 @@ export * from './models/IPbxEventsModel'; export * from './models/IPushTokenModel'; export * from './models/IPermissionsModel'; export * from './models/IReadReceiptsModel'; +export * from './models/IMessageReadsModel'; export * from './models/IReportsModel'; export * from './models/IRolesModel'; export * from './models/IRoomsModel'; diff --git a/packages/model-typings/src/models/IMessageReadsModel.ts b/packages/model-typings/src/models/IMessageReadsModel.ts new file mode 100644 index 0000000000000..fd696d8403feb --- /dev/null +++ b/packages/model-typings/src/models/IMessageReadsModel.ts @@ -0,0 +1,12 @@ +import type { UpdateResult } from 'mongodb'; +import type { MessageReads, IUser, IMessage } from '@rocket.chat/core-typings'; + +import type { IBaseModel } from './IBaseModel'; + +export interface IMessageReadsModel extends IBaseModel { + findOneByUserIdAndThreadId(userId: IUser['_id'], tmid: IMessage['_id']): Promise; + getMinimumLastSeenByThreadId(tmid: string): Promise; + updateReadTimestampByUserIdAndThreadId(userId: IUser['_id'], tmid: IMessage['_id']): Promise; + countByThreadAndUserIds(tmid: IMessage['_id'], userIds: IUser['_id'][]): Promise; + countByThreadId(tmid: IMessage['_id']): Promise; +} diff --git a/packages/model-typings/src/models/ISubscriptionsModel.ts b/packages/model-typings/src/models/ISubscriptionsModel.ts index f0b80986dfec6..ac2a6c5975bb3 100644 --- a/packages/model-typings/src/models/ISubscriptionsModel.ts +++ b/packages/model-typings/src/models/ISubscriptionsModel.ts @@ -12,12 +12,16 @@ export interface ISubscriptionsModel extends IBaseModel { findByRoomId(roomId: string, options?: FindOptions): FindCursor; + findUnarchivedByRoomId(roomId: string, options?: FindOptions): FindCursor; + findByRoomIdAndNotUserId(roomId: string, userId: string, options?: FindOptions): FindCursor; findByLivechatRoomIdAndNotUserId(roomId: string, userId: string, options?: FindOptions): FindCursor; countByRoomIdAndUserId(rid: string, uid: string | undefined): Promise; + countUnarchivedByRoomId(rid: string): Promise; + isUserInRole(uid: IUser['_id'], roleId: IRole['_id'], rid?: IRoom['_id']): Promise; setAsReadByRoomIdAndUserId( diff --git a/packages/models/src/index.ts b/packages/models/src/index.ts index 64d16afba0db8..e2c2969912392 100644 --- a/packages/models/src/index.ts +++ b/packages/models/src/index.ts @@ -42,6 +42,7 @@ import type { IPushTokenModel, IPermissionsModel, IReadReceiptsModel, + IMessageReadsModel, IReportsModel, IRolesModel, IRoomsModel, @@ -116,6 +117,7 @@ export const PbxEvents = proxify('IPbxEventsModel'); export const PushToken = proxify('IPushTokenModel'); export const Permissions = proxify('IPermissionsModel'); export const ReadReceipts = proxify('IReadReceiptsModel'); +export const MessageReads = proxify('IMessageReadsModel'); export const Reports = proxify('IReportsModel'); export const Roles = proxify('IRolesModel'); export const Rooms = proxify('IRoomsModel');