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
This commit is contained in:
Philip Okugbe
2026-09-30 01:10:16 +01:00
committed by GitHub
parent 870b71d2aa
commit 01a139c3d1
10 changed files with 113 additions and 23 deletions
@@ -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),
},
}),
@@ -146,8 +146,14 @@ export class CollaborationGateway {
eventName: TName,
documentName: string,
payload: Parameters<CollabEventHandlers[TName]>[1],
onlyIfOpen = false,
) {
return this.redisSync?.handleEvent(eventName, documentName, payload);
return this.redisSync?.handleEvent(
eventName,
documentName,
payload,
onlyIfOpen,
);
}
openDirectConnection(documentName: string, context?: any) {
@@ -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<string>();
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);
@@ -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);
}
@@ -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());
}
@@ -235,13 +235,11 @@ export class RedisSyncExtension<TCE extends CustomEvents> 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);
}
@@ -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, () => {
@@ -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,
@@ -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<string, Set<string>>();