From 01a139c3d1fb9c2f88190821933043557cddd22b Mon Sep 17 00:00:00 2001 From: Philip Okugbe <16838612+Philipinho@users.noreply.github.com> Date: Wed, 30 Sep 2026 01:10:16 +0100 Subject: [PATCH] fix: skip page update notifications for users viewing the page (#2530) * fix: drop malformed awareness before broadcast * clear interval * pass userId and avatar to awareness * fix: skip page update notifications for editors and users viewing the page --- .../features/editor/extensions/extensions.ts | 2 + .../collaboration/collaboration.gateway.ts | 8 ++- .../collaboration/collaboration.handler.ts | 11 ++++ .../src/collaboration/collaboration.util.ts | 4 ++ .../extensions/persistence.extension.ts | 53 +++++++++++++++++-- .../redis-sync/redis-sync.extension.ts | 10 ++-- .../src/collaboration/server/collab-main.ts | 8 +++ .../core/notification/notification.module.ts | 3 +- .../services/page.notification.ts | 35 ++++++++---- apps/server/src/ee | 2 +- 10 files changed, 113 insertions(+), 23 deletions(-) diff --git a/apps/client/src/features/editor/extensions/extensions.ts b/apps/client/src/features/editor/extensions/extensions.ts index 7c685ac65..76a4cf88f 100644 --- a/apps/client/src/features/editor/extensions/extensions.ts +++ b/apps/client/src/features/editor/extensions/extensions.ts @@ -471,7 +471,9 @@ export const collabExtensions: CollabExtensions = (provider, user) => [ CollaborationCaret.configure({ provider, user: { + id: user.id, name: user.name, + avatarUrl: user.avatarUrl, color: randomElement(userColors), }, }), diff --git a/apps/server/src/collaboration/collaboration.gateway.ts b/apps/server/src/collaboration/collaboration.gateway.ts index 254b8d430..5556fef19 100644 --- a/apps/server/src/collaboration/collaboration.gateway.ts +++ b/apps/server/src/collaboration/collaboration.gateway.ts @@ -146,8 +146,14 @@ export class CollaborationGateway { eventName: TName, documentName: string, payload: Parameters[1], + onlyIfOpen = false, ) { - return this.redisSync?.handleEvent(eventName, documentName, payload); + return this.redisSync?.handleEvent( + eventName, + documentName, + payload, + onlyIfOpen, + ); } openDirectConnection(documentName: string, context?: any) { diff --git a/apps/server/src/collaboration/collaboration.handler.ts b/apps/server/src/collaboration/collaboration.handler.ts index 992f9b74b..2df9d9721 100644 --- a/apps/server/src/collaboration/collaboration.handler.ts +++ b/apps/server/src/collaboration/collaboration.handler.ts @@ -21,6 +21,17 @@ export class CollaborationHandler { getHandlers(hocuspocus: Hocuspocus) { return { + getConnectedUserIds: async (documentName: string) => { + const document = hocuspocus.documents.get(documentName); + if (!document) return []; + + const userIds = new Set(); + for (const state of document.awareness.getStates().values()) { + const userId = state?.user?.id; + if (typeof userId === 'string') userIds.add(userId); + } + return [...userIds]; + }, alterState: async (documentName: string, payload: { pageId: string }) => { // dummy // this.logger.log('Processing', documentName, payload); diff --git a/apps/server/src/collaboration/collaboration.util.ts b/apps/server/src/collaboration/collaboration.util.ts index e5c55c77a..1c527ea2a 100644 --- a/apps/server/src/collaboration/collaboration.util.ts +++ b/apps/server/src/collaboration/collaboration.util.ts @@ -247,3 +247,7 @@ export function jsonToMarkdown(tiptapJson: any): string { const html = jsonToHtml(tiptapJson); return htmlToMarkdown(html); } + +export function isRenderableObject(value: unknown): boolean { + return typeof value === 'object' && value !== null && !Array.isArray(value); +} diff --git a/apps/server/src/collaboration/extensions/persistence.extension.ts b/apps/server/src/collaboration/extensions/persistence.extension.ts index 85efb5ec4..f13ff4407 100644 --- a/apps/server/src/collaboration/extensions/persistence.extension.ts +++ b/apps/server/src/collaboration/extensions/persistence.extension.ts @@ -1,5 +1,6 @@ import { afterUnloadDocumentPayload, + beforeHandleAwarenessPayload, Extension, onChangePayload, onLoadDocumentPayload, @@ -8,7 +9,12 @@ import { import * as Y from 'yjs'; import { Injectable, Logger } from '@nestjs/common'; import { TiptapTransformer } from '@hocuspocus/transformer'; -import { getPageId, jsonToText, tiptapExtensions } from '../collaboration.util'; +import { + getPageId, + isRenderableObject, + jsonToText, + tiptapExtensions, +} from '../collaboration.util'; import { PageRepo } from '@docmost/db/repos/page/page.repo'; import { InjectKysely } from 'nestjs-kysely'; import { KyselyDB } from '@docmost/db/types/kysely.types'; @@ -132,6 +138,8 @@ export class PersistenceExtension implements Extension { return; } + await this.collabHistory.addContributors(pageId, editingUserIds); + let contributorIds = undefined; try { const existingContributors = page.contributorIds || []; @@ -162,6 +170,10 @@ export class PersistenceExtension implements Extension { }); } catch (err) { this.logger.error(`Failed to update page ${pageId}`, err); + page = null; + editingUserIds.forEach((userId) => + this.trackContributor(documentName, userId), + ); } if (page) { @@ -184,8 +196,6 @@ export class PersistenceExtension implements Extension { } if (page) { - await this.collabHistory.addContributors(pageId, editingUserIds); - const mentions = extractMentions(tiptapJson); const userMentions = extractUserMentions(mentions); @@ -217,12 +227,45 @@ export class PersistenceExtension implements Extension { } } + // Drop malformed awareness before it is broadcast + async beforeHandleAwareness({ + states, + context, + }: beforeHandleAwarenessPayload) { + const user = context?.user; + + for (const [clientId, state] of states) { + if (!isRenderableObject(state)) { + states.delete(clientId); + continue; + } + + if ('user' in state && !isRenderableObject(state.user)) { + delete state.user; + } + + if (state.user && user) { + state.user.id = user.id; + state.user.avatarUrl = user.avatarUrl ?? null; + } + + if ( + 'cursor' in state && + state.cursor !== null && + !isRenderableObject(state.cursor) + ) { + delete state.cursor; + } + } + } + async onChange(data: onChangePayload) { - const documentName = data.documentName; const userId = data.context?.user?.id; - if (!userId) return; + this.trackContributor(data.documentName, userId); + } + private trackContributor(documentName: string, userId: string) { if (!this.contributors.has(documentName)) { this.contributors.set(documentName, new Set()); } diff --git a/apps/server/src/collaboration/extensions/redis-sync/redis-sync.extension.ts b/apps/server/src/collaboration/extensions/redis-sync/redis-sync.extension.ts index 42139097f..3a667f01f 100644 --- a/apps/server/src/collaboration/extensions/redis-sync/redis-sync.extension.ts +++ b/apps/server/src/collaboration/extensions/redis-sync/redis-sync.extension.ts @@ -235,13 +235,11 @@ export class RedisSyncExtension implements Extension { }; async maintainLock(documentName: string) { + clearInterval(this.locks[documentName]); this.locks[documentName] = setInterval(() => { - this.pub.set( - this.getKey(documentName), - this.serverId, - 'PX', - this.lockTTL, - ); + this.pub + .set(this.getKey(documentName), this.serverId, 'PX', this.lockTTL) + .catch(() => {}); }, this.lockTTL / 2); } diff --git a/apps/server/src/collaboration/server/collab-main.ts b/apps/server/src/collaboration/server/collab-main.ts index 3b3de2430..89b1605b5 100644 --- a/apps/server/src/collaboration/server/collab-main.ts +++ b/apps/server/src/collaboration/server/collab-main.ts @@ -37,6 +37,14 @@ async function bootstrap() { const logger = new Logger('CollabServer'); + process.on('unhandledRejection', (reason, promise) => { + logger.error(`UnhandledRejection, reason: ${reason}`, promise); + }); + + process.on('uncaughtException', (error) => { + logger.error('UncaughtException:', error); + }); + const port = process.env.COLLAB_PORT || 3001; const host = process.env.HOST || '0.0.0.0'; await app.listen(port, host, () => { diff --git a/apps/server/src/core/notification/notification.module.ts b/apps/server/src/core/notification/notification.module.ts index 9aa452a5c..543ca0cc3 100644 --- a/apps/server/src/core/notification/notification.module.ts +++ b/apps/server/src/core/notification/notification.module.ts @@ -6,9 +6,10 @@ import { CommentNotificationService } from './services/comment.notification'; import { PageNotificationService } from './services/page.notification'; import { VerificationNotificationService } from './services/verification.notification'; import { PageUpdateEmailRateLimiter } from './services/page-update-email-rate-limiter'; +import { CollaborationModule } from '../../collaboration/collaboration.module'; @Module({ - imports: [], + imports: [CollaborationModule], controllers: [NotificationController], providers: [ NotificationService, diff --git a/apps/server/src/core/notification/services/page.notification.ts b/apps/server/src/core/notification/services/page.notification.ts index 77ab967aa..dcbc6e865 100644 --- a/apps/server/src/core/notification/services/page.notification.ts +++ b/apps/server/src/core/notification/services/page.notification.ts @@ -21,6 +21,7 @@ import { PageUpdateDigestEmail } from '@docmost/transactional/emails/page-update import { PermissionGrantedEmail } from '@docmost/transactional/emails/permission-granted-email'; import { getPageTitle } from '../../../common/helpers'; import { QueueJob, QueueName } from '../../../integrations/queue/constants'; +import { CollaborationGateway } from '../../../collaboration/collaboration.gateway'; const PAGE_UPDATE_COOLDOWN_HOURS = 7; const DIGEST_DELAY_MS = 12 * 60 * 60 * 1000; // 12 hours @@ -37,6 +38,7 @@ export class PageNotificationService { private readonly pagePermissionRepo: PagePermissionRepo, private readonly watcherRepo: WatcherRepo, private readonly rateLimiter: PageUpdateEmailRateLimiter, + private readonly collaborationGateway: CollaborationGateway, @InjectQueue(QueueName.NOTIFICATION_QUEUE) private notificationQueue: Queue, ) {} @@ -186,8 +188,22 @@ export class PageNotificationService { if (watcherIds.length === 0) return; - const actorSet = new Set(actorIds); - const candidateIds = watcherIds.filter((id) => !actorSet.has(id)); + let connectedUserIds: string[] = []; + try { + connectedUserIds = + (await this.collaborationGateway.handleYjsEvent( + 'getConnectedUserIds', + `page.${pageId}`, + undefined, + true, + )) ?? []; + } catch (err) { + this.logger.warn( + `Failed to get connected users for page ${pageId}: ${err?.['message']}`, + ); + } + const excludedIds = new Set([...actorIds, ...connectedUserIds]); + const candidateIds = watcherIds.filter((id) => !excludedIds.has(id)); if (candidateIds.length === 0) return; const eligibleUsers = await this.getEligiblePageUpdateUsers(candidateIds); @@ -370,13 +386,14 @@ export class PageNotificationService { const pages = spaceFilteredPages.filter((p) => accessiblePageIds.has(p.id)); if (pages.length === 0) return; - const actors = actorIds.length > 0 - ? await this.db - .selectFrom('users') - .select(['id', 'name']) - .where('id', 'in', actorIds) - .execute() - : []; + const actors = + actorIds.length > 0 + ? await this.db + .selectFrom('users') + .select(['id', 'name']) + .where('id', 'in', actorIds) + .execute() + : []; const actorMap = new Map(actors.map((a) => [a.id, a.name])); const pageActors = new Map>(); diff --git a/apps/server/src/ee b/apps/server/src/ee index 9134679e3..8190e2182 160000 --- a/apps/server/src/ee +++ b/apps/server/src/ee @@ -1 +1 @@ -Subproject commit 9134679e37e8ccd0f79e6f6c05da8e238778420d +Subproject commit 8190e2182cbbcc5e291217a13147fcd246ede3fe