From a444941b793c78a09739a30c0a6204119c7ceb1b Mon Sep 17 00:00:00 2001 From: DarkSky <25152247+darkskygit@users.noreply.github.com> Date: Tue, 15 Jul 2025 20:21:42 +0800 Subject: [PATCH] fix(server): delay send mail if retry many times (#13225) fix AF-2748 ## Summary by CodeRabbit * **New Features** * Improved mail sending job with adaptive retry delays based on elapsed time, enhancing reliability of email delivery. * **Chores** * Updated job payload to include a start time for better retry management. * Added an internal delay utility to support asynchronous pause in processes. --- .../backend/server/src/base/utils/promise.ts | 6 ++++++ packages/backend/server/src/core/mail/job.ts | 20 +++++++++++++++---- .../backend/server/src/core/mail/mailer.ts | 14 ++++++++++--- 3 files changed, 33 insertions(+), 7 deletions(-) diff --git a/packages/backend/server/src/base/utils/promise.ts b/packages/backend/server/src/base/utils/promise.ts index f59f1e6488..27023dc341 100644 --- a/packages/backend/server/src/base/utils/promise.ts +++ b/packages/backend/server/src/base/utils/promise.ts @@ -1,3 +1,5 @@ +import { setTimeout } from 'node:timers/promises'; + import { defer as rxjsDefer, retry } from 'rxjs'; export class RetryablePromise extends Promise { @@ -48,3 +50,7 @@ export function defer(dispose: () => Promise) { [Symbol.asyncDispose]: dispose, }; } + +export function sleep(ms: number): Promise { + return setTimeout(ms); +} diff --git a/packages/backend/server/src/core/mail/job.ts b/packages/backend/server/src/core/mail/job.ts index b2b7ef3da1..fdea182419 100644 --- a/packages/backend/server/src/core/mail/job.ts +++ b/packages/backend/server/src/core/mail/job.ts @@ -1,7 +1,7 @@ import { Injectable } from '@nestjs/common'; import { getStreamAsBuffer } from 'get-stream'; -import { JOB_SIGNAL, OnJob } from '../../base'; +import { JOB_SIGNAL, OnJob, sleep } from '../../base'; import { type MailName, MailProps, Renderers } from '../../mails'; import { UserProps, WorkspaceProps } from '../../mails/components'; import { Models } from '../../models'; @@ -34,7 +34,7 @@ type SendMailJob> = { declare global { interface Jobs { - 'notification.sendMail': { + 'notification.sendMail': { startTime: number } & { [K in MailName]: SendMailJob; }[MailName]; } @@ -50,7 +50,12 @@ export class MailJob { ) {} @OnJob('notification.sendMail') - async sendMail({ name, to, props }: Jobs['notification.sendMail']) { + async sendMail({ + startTime, + name, + to, + props, + }: Jobs['notification.sendMail']) { let options: Partial = {}; for (const key in props) { @@ -100,8 +105,15 @@ export class MailJob { )), ...options, }); + if (result === false) { + // wait for a while before retrying + const elapsed = Date.now() - startTime; + const retryDelay = Math.min(30 * 1000, Math.round(elapsed / 2000) * 1000); + await sleep(retryDelay); + return JOB_SIGNAL.Retry; + } - return result === false ? JOB_SIGNAL.Retry : undefined; + return undefined; } private async fetchWorkspaceProps(workspaceId: string) { diff --git a/packages/backend/server/src/core/mail/mailer.ts b/packages/backend/server/src/core/mail/mailer.ts index c1d459342c..4ad492e23c 100644 --- a/packages/backend/server/src/core/mail/mailer.ts +++ b/packages/backend/server/src/core/mail/mailer.ts @@ -15,11 +15,14 @@ export class Mailer { * * @note never throw */ - async trySend(command: Jobs['notification.sendMail']) { + async trySend(command: Omit) { return this.send(command, true); } - async send(command: Jobs['notification.sendMail'], suppressError = false) { + async send( + command: Omit, + suppressError = false + ) { if (!this.sender.configured) { if (suppressError) { return false; @@ -28,7 +31,12 @@ export class Mailer { } try { - await this.queue.add('notification.sendMail', command); + await this.queue.add( + 'notification.sendMail', + Object.assign({}, command, { + startTime: Date.now(), + }) as Jobs['notification.sendMail'] + ); return true; } catch { return false;