mirror of
https://github.com/docmost/docmost.git
synced 2026-08-30 02:25:01 +08:00
digests
This commit is contained in:
@@ -4,6 +4,7 @@ import { NotificationController } from './notification.controller';
|
|||||||
import { NotificationProcessor } from './notification.processor';
|
import { NotificationProcessor } from './notification.processor';
|
||||||
import { CommentNotificationService } from './services/comment.notification';
|
import { CommentNotificationService } from './services/comment.notification';
|
||||||
import { PageNotificationService } from './services/page.notification';
|
import { PageNotificationService } from './services/page.notification';
|
||||||
|
import { PageUpdateEmailRateLimiter } from './services/page-update-email-rate-limiter';
|
||||||
|
|
||||||
@Module({
|
@Module({
|
||||||
imports: [],
|
imports: [],
|
||||||
@@ -13,6 +14,7 @@ import { PageNotificationService } from './services/page.notification';
|
|||||||
NotificationProcessor,
|
NotificationProcessor,
|
||||||
CommentNotificationService,
|
CommentNotificationService,
|
||||||
PageNotificationService,
|
PageNotificationService,
|
||||||
|
PageUpdateEmailRateLimiter,
|
||||||
],
|
],
|
||||||
exports: [NotificationService],
|
exports: [NotificationService],
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -87,6 +87,12 @@ export class NotificationProcessor
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
case QueueJob.PAGE_UPDATE_DIGEST: {
|
||||||
|
const { userId } = job.data as unknown as { userId: string };
|
||||||
|
await this.pageNotificationService.processDigest(userId, appUrl);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
default:
|
default:
|
||||||
this.logger.warn(`Unknown notification job: ${job.name}`);
|
this.logger.warn(`Unknown notification job: ${job.name}`);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,44 @@
|
|||||||
|
import { Injectable } from '@nestjs/common';
|
||||||
|
import { RedisService } from '@nestjs-labs/nestjs-ioredis';
|
||||||
|
import type { Redis } from 'ioredis';
|
||||||
|
|
||||||
|
const KEY_PREFIX = 'page-update:emails:';
|
||||||
|
const DIGEST_PREFIX = 'page-update:digest:';
|
||||||
|
const TTL_SECONDS = 86400; // 24 hours
|
||||||
|
const MAX_IMMEDIATE_EMAILS = 10;
|
||||||
|
|
||||||
|
@Injectable()
|
||||||
|
export class PageUpdateEmailRateLimiter {
|
||||||
|
private readonly redis: Redis;
|
||||||
|
|
||||||
|
constructor(private readonly redisService: RedisService) {
|
||||||
|
this.redis = this.redisService.getOrThrow();
|
||||||
|
}
|
||||||
|
|
||||||
|
async canSendEmail(userId: string): Promise<boolean> {
|
||||||
|
const key = KEY_PREFIX + userId;
|
||||||
|
const count = await this.redis.incr(key);
|
||||||
|
await this.redis.expire(key, TTL_SECONDS, 'NX');
|
||||||
|
return count <= MAX_IMMEDIATE_EMAILS;
|
||||||
|
}
|
||||||
|
|
||||||
|
async addToDigest(userId: string, notificationId: string): Promise<boolean> {
|
||||||
|
const key = DIGEST_PREFIX + userId;
|
||||||
|
const isNew = (await this.redis.llen(key)) === 0;
|
||||||
|
await this.redis.rpush(key, notificationId);
|
||||||
|
await this.redis.expire(key, TTL_SECONDS);
|
||||||
|
return isNew;
|
||||||
|
}
|
||||||
|
|
||||||
|
async popDigest(userId: string): Promise<string[]> {
|
||||||
|
const key = DIGEST_PREFIX + userId;
|
||||||
|
const [ids] = await this.redis
|
||||||
|
.multi()
|
||||||
|
.lrange(key, 0, -1)
|
||||||
|
.del(key)
|
||||||
|
.exec();
|
||||||
|
|
||||||
|
return (ids?.[1] as string[]) ?? [];
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
@@ -1,5 +1,7 @@
|
|||||||
import { Injectable } from '@nestjs/common';
|
import { Injectable, Logger } from '@nestjs/common';
|
||||||
import { InjectKysely } from 'nestjs-kysely';
|
import { InjectKysely } from 'nestjs-kysely';
|
||||||
|
import { InjectQueue } from '@nestjs/bullmq';
|
||||||
|
import { Queue } from 'bullmq';
|
||||||
import { KyselyDB } from '@docmost/db/types/kysely.types';
|
import { KyselyDB } from '@docmost/db/types/kysely.types';
|
||||||
import {
|
import {
|
||||||
IPageMentionNotificationJob,
|
IPageMentionNotificationJob,
|
||||||
@@ -12,15 +14,21 @@ import { NotificationRepo } from '@docmost/db/repos/notification/notification.re
|
|||||||
import { SpaceMemberRepo } from '@docmost/db/repos/space/space-member.repo';
|
import { SpaceMemberRepo } from '@docmost/db/repos/space/space-member.repo';
|
||||||
import { PagePermissionRepo } from '@docmost/db/repos/page/page-permission.repo';
|
import { PagePermissionRepo } from '@docmost/db/repos/page/page-permission.repo';
|
||||||
import { WatcherRepo } from '@docmost/db/repos/watcher/watcher.repo';
|
import { WatcherRepo } from '@docmost/db/repos/watcher/watcher.repo';
|
||||||
|
import { PageUpdateEmailRateLimiter } from './page-update-email-rate-limiter';
|
||||||
import { PageMentionEmail } from '@docmost/transactional/emails/page-mention-email';
|
import { PageMentionEmail } from '@docmost/transactional/emails/page-mention-email';
|
||||||
import { PageUpdateEmail } from '@docmost/transactional/emails/page-update-email';
|
import { PageUpdateEmail } from '@docmost/transactional/emails/page-update-email';
|
||||||
|
import { PageUpdateDigestEmail } from '@docmost/transactional/emails/page-update-digest-email';
|
||||||
import { PermissionGrantedEmail } from '@docmost/transactional/emails/permission-granted-email';
|
import { PermissionGrantedEmail } from '@docmost/transactional/emails/permission-granted-email';
|
||||||
import { getPageTitle } from '../../../common/helpers';
|
import { getPageTitle } from '../../../common/helpers';
|
||||||
|
import { QueueJob, QueueName } from '../../../integrations/queue/constants';
|
||||||
|
|
||||||
const PAGE_UPDATE_COOLDOWN_HOURS = 7;
|
const PAGE_UPDATE_COOLDOWN_HOURS = 7;
|
||||||
|
const DIGEST_DELAY_MS = 3 * 60 * 60 * 1000; // 3 hours
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class PageNotificationService {
|
export class PageNotificationService {
|
||||||
|
private readonly logger = new Logger(PageNotificationService.name);
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
@InjectKysely() private readonly db: KyselyDB,
|
@InjectKysely() private readonly db: KyselyDB,
|
||||||
private readonly notificationService: NotificationService,
|
private readonly notificationService: NotificationService,
|
||||||
@@ -28,6 +36,8 @@ export class PageNotificationService {
|
|||||||
private readonly spaceMemberRepo: SpaceMemberRepo,
|
private readonly spaceMemberRepo: SpaceMemberRepo,
|
||||||
private readonly pagePermissionRepo: PagePermissionRepo,
|
private readonly pagePermissionRepo: PagePermissionRepo,
|
||||||
private readonly watcherRepo: WatcherRepo,
|
private readonly watcherRepo: WatcherRepo,
|
||||||
|
private readonly rateLimiter: PageUpdateEmailRateLimiter,
|
||||||
|
@InjectQueue(QueueName.NOTIFICATION_QUEUE) private notificationQueue: Queue,
|
||||||
) {}
|
) {}
|
||||||
|
|
||||||
async processPageMention(data: IPageMentionNotificationJob, appUrl: string) {
|
async processPageMention(data: IPageMentionNotificationJob, appUrl: string) {
|
||||||
@@ -220,17 +230,28 @@ export class PageNotificationService {
|
|||||||
});
|
});
|
||||||
if (!notification) continue;
|
if (!notification) continue;
|
||||||
|
|
||||||
await this.notificationService.queueEmail(
|
const canSend = await this.rateLimiter.canSendEmail(userId);
|
||||||
userId,
|
if (canSend) {
|
||||||
notification.id,
|
await this.notificationService.queueEmail(
|
||||||
`${actor.name} updated ${pageTitle}`,
|
userId,
|
||||||
PageUpdateEmail({
|
notification.id,
|
||||||
actorName: actor.name,
|
`${actor.name} updated ${pageTitle}`,
|
||||||
pageTitle,
|
PageUpdateEmail({
|
||||||
pageUrl: basePageUrl,
|
actorName: actor.name,
|
||||||
}),
|
pageTitle,
|
||||||
NotificationType.PAGE_UPDATED,
|
pageUrl: basePageUrl,
|
||||||
);
|
}),
|
||||||
|
NotificationType.PAGE_UPDATED,
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
const isFirst = await this.rateLimiter.addToDigest(
|
||||||
|
userId,
|
||||||
|
notification.id,
|
||||||
|
);
|
||||||
|
if (isFirst) {
|
||||||
|
await this.scheduleDigest(userId, workspaceId);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -255,6 +276,66 @@ export class PageNotificationService {
|
|||||||
.map((u) => u.id);
|
.map((u) => u.id);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private async scheduleDigest(
|
||||||
|
userId: string,
|
||||||
|
workspaceId: string,
|
||||||
|
): Promise<void> {
|
||||||
|
const jobId = `page-update-digest:${userId}`;
|
||||||
|
await this.notificationQueue
|
||||||
|
.add(
|
||||||
|
QueueJob.PAGE_UPDATE_DIGEST,
|
||||||
|
{ userId, workspaceId },
|
||||||
|
{ jobId, delay: DIGEST_DELAY_MS },
|
||||||
|
)
|
||||||
|
.catch((err) => {
|
||||||
|
this.logger.error(
|
||||||
|
`Failed to schedule digest for ${userId}: ${err.message}`,
|
||||||
|
);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async processDigest(userId: string, appUrl: string): Promise<void> {
|
||||||
|
const notificationIds = await this.rateLimiter.popDigest(userId);
|
||||||
|
if (notificationIds.length === 0) return;
|
||||||
|
|
||||||
|
const notifications = await this.db
|
||||||
|
.selectFrom('notifications')
|
||||||
|
.select(['id', 'pageId'])
|
||||||
|
.where('id', 'in', notificationIds)
|
||||||
|
.execute();
|
||||||
|
|
||||||
|
if (notifications.length === 0) return;
|
||||||
|
|
||||||
|
const pageIds = [...new Set(notifications.map((n) => n.pageId).filter(Boolean))];
|
||||||
|
|
||||||
|
const pages = await this.db
|
||||||
|
.selectFrom('pages')
|
||||||
|
.innerJoin('spaces', 'spaces.id', 'pages.spaceId')
|
||||||
|
.select([
|
||||||
|
'pages.id',
|
||||||
|
'pages.title',
|
||||||
|
'pages.slugId',
|
||||||
|
'spaces.slug as spaceSlug',
|
||||||
|
])
|
||||||
|
.where('pages.id', 'in', pageIds)
|
||||||
|
.execute();
|
||||||
|
|
||||||
|
const pageUpdates = pages.map((p) => ({
|
||||||
|
title: getPageTitle(p.title),
|
||||||
|
url: `${appUrl}/s/${p.spaceSlug}/p/${p.slugId}`,
|
||||||
|
}));
|
||||||
|
|
||||||
|
if (pageUpdates.length === 0) return;
|
||||||
|
|
||||||
|
await this.notificationService.queueEmail(
|
||||||
|
userId,
|
||||||
|
notificationIds[0],
|
||||||
|
`${pageUpdates.length} pages were updated`,
|
||||||
|
PageUpdateDigestEmail({ pageUpdates }),
|
||||||
|
NotificationType.PAGE_UPDATED,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
private async getPageContext(
|
private async getPageContext(
|
||||||
actorId: string,
|
actorId: string,
|
||||||
pageId: string,
|
pageId: string,
|
||||||
|
|||||||
+1
-1
Submodule apps/server/src/ee updated: f486726088...05f1c816a8
@@ -69,6 +69,7 @@ export enum QueueJob {
|
|||||||
COMMENT_RESOLVED_NOTIFICATION = 'comment-resolved-notification',
|
COMMENT_RESOLVED_NOTIFICATION = 'comment-resolved-notification',
|
||||||
PAGE_MENTION_NOTIFICATION = 'page-mention-notification',
|
PAGE_MENTION_NOTIFICATION = 'page-mention-notification',
|
||||||
PAGE_PERMISSION_GRANTED = 'page-permission-granted',
|
PAGE_PERMISSION_GRANTED = 'page-permission-granted',
|
||||||
|
PAGE_UPDATE_DIGEST = 'page-update-digest',
|
||||||
|
|
||||||
AUDIT_LOG = 'audit-log',
|
AUDIT_LOG = 'audit-log',
|
||||||
AUDIT_CLEANUP = 'audit-cleanup',
|
AUDIT_CLEANUP = 'audit-cleanup',
|
||||||
|
|||||||
@@ -0,0 +1,42 @@
|
|||||||
|
import { Link, Section, Text } from '@react-email/components';
|
||||||
|
import * as React from 'react';
|
||||||
|
import { content, link, paragraph } from '../css/styles';
|
||||||
|
import { MailBody } from '../partials/partials';
|
||||||
|
|
||||||
|
interface PageUpdate {
|
||||||
|
title: string;
|
||||||
|
url: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
interface Props {
|
||||||
|
pageUpdates: PageUpdate[];
|
||||||
|
}
|
||||||
|
|
||||||
|
export const PageUpdateDigestEmail = ({ pageUpdates }: Props) => {
|
||||||
|
return (
|
||||||
|
<MailBody>
|
||||||
|
<Section style={content}>
|
||||||
|
<Text style={paragraph}>Hi there,</Text>
|
||||||
|
<Text style={paragraph}>
|
||||||
|
The following {pageUpdates.length} pages you watch were updated:
|
||||||
|
</Text>
|
||||||
|
{pageUpdates.map((page, i) => (
|
||||||
|
<Text key={i} style={listItem}>
|
||||||
|
{'• '}
|
||||||
|
<Link href={page.url} style={link}>
|
||||||
|
{page.title}
|
||||||
|
</Link>
|
||||||
|
</Text>
|
||||||
|
))}
|
||||||
|
</Section>
|
||||||
|
</MailBody>
|
||||||
|
);
|
||||||
|
};
|
||||||
|
|
||||||
|
const listItem = {
|
||||||
|
...paragraph,
|
||||||
|
margin: '4px 0',
|
||||||
|
lineHeight: 1.4,
|
||||||
|
};
|
||||||
|
|
||||||
|
export default PageUpdateDigestEmail;
|
||||||
Reference in New Issue
Block a user