fix: skip page update notifications for editors and users viewing the page

This commit is contained in:
Philipinho
2026-09-30 01:06:57 +01:00
parent f2d81744bb
commit 857edb90fe
6 changed files with 63 additions and 15 deletions
@@ -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);
@@ -138,6 +138,8 @@ export class PersistenceExtension implements Extension {
return;
}
await this.collabHistory.addContributors(pageId, editingUserIds);
let contributorIds = undefined;
try {
const existingContributors = page.contributorIds || [];
@@ -168,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) {
@@ -190,8 +196,6 @@ export class PersistenceExtension implements Extension {
}
if (page) {
await this.collabHistory.addContributors(pageId, editingUserIds);
const mentions = extractMentions(tiptapJson);
const userMentions = extractUserMentions(mentions);
@@ -256,11 +260,12 @@ export class PersistenceExtension implements Extension {
}
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());
}
@@ -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>>();