From b9ce7373a47500da972ca5d4766724df1ff4abb1 Mon Sep 17 00:00:00 2001 From: loks0n <22452787+loks0n@users.noreply.github.com> Date: Fri, 13 Mar 2026 19:20:56 +0000 Subject: [PATCH] Refactor mails queue to stateless publisher --- app/controllers/api/account.php | 273 ++++++++++-------- app/controllers/api/projects.php | 42 +-- app/controllers/shared/api.php | 8 +- app/init/resources.php | 9 +- app/worker.php | 9 +- src/Appwrite/Event/Message/Mail.php | 61 ++++ src/Appwrite/Event/Publisher/Mail.php | 27 ++ .../Http/Account/MFA/Challenges/Create.php | 54 ++-- .../Health/Http/Health/Queue/Failed/Get.php | 8 +- .../Health/Http/Health/Queue/Mails/Get.php | 8 +- .../Modules/Teams/Http/Memberships/Create.php | 47 +-- .../Platform/Workers/Certificates.php | 45 +-- src/Appwrite/Platform/Workers/Migrations.php | 54 ++-- src/Appwrite/Platform/Workers/Webhooks.php | 45 ++- 14 files changed, 413 insertions(+), 277 deletions(-) create mode 100644 src/Appwrite/Event/Message/Mail.php create mode 100644 src/Appwrite/Event/Publisher/Mail.php diff --git a/app/controllers/api/account.php b/app/controllers/api/account.php index a780bfdac3..03a38ff909 100644 --- a/app/controllers/api/account.php +++ b/app/controllers/api/account.php @@ -12,8 +12,9 @@ use Appwrite\Auth\Validator\Phone; use Appwrite\Detector\Detector; use Appwrite\Event\Delete; use Appwrite\Event\Event; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; use Appwrite\Event\Messaging; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Extend\Exception; use Appwrite\Hooks\Hooks; use Appwrite\Network\Validator\Email as EmailValidator; @@ -75,7 +76,7 @@ use Utopia\Validator\WhiteList; $oauthDefaultSuccess = '/console/auth/oauth2/success'; $oauthDefaultFailure = '/console/auth/oauth2/failure'; -function sendSessionAlert(Locale $locale, Document $user, Document $project, array $platform, Document $session, Mail $queueForMails) +function sendSessionAlert(Locale $locale, Document $user, Document $project, array $platform, Document $session, MailsPublisher $publisherForMails) { $subject = $locale->getText("emails.sessionAlert.subject"); $preview = $locale->getText("emails.sessionAlert.preview"); @@ -109,6 +110,9 @@ function sendSessionAlert(Locale $locale, Document $user, Document $project, arr $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -120,12 +124,13 @@ function sendSessionAlert(Locale $locale, Document $user, Document $project, arr $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -142,10 +147,9 @@ function sendSessionAlert(Locale $locale, Document $user, Document $project, arr $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } // session alerts should always have a client name! @@ -190,24 +194,26 @@ function sendSessionAlert(Locale $locale, Document $user, Document $project, arr $email = $user->getAttribute('email'); - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->setBodyTemplate($bodyTemplate) - ->appendVariables($emailVariables) - ->setRecipient($email); - - // since this is console project, set email sender name! if ($smtpBaseTemplate === APP_BRANDED_EMAIL_BASE_TEMPLATE) { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $email, + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + bodyTemplate: $bodyTemplate, + customMailOptions: $customMailOptions, + platform: $platform, + )); } -$createSession = function (string $userId, string $secret, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, Reader $geodb, Event $queueForEvents, Mail $queueForMails, Store $store, ProofsToken $proofForToken, ProofsCode $proofForCode, Authorization $authorization) { +$createSession = function (string $userId, string $secret, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, Reader $geodb, Event $queueForEvents, MailsPublisher $publisherForMails, Store $store, ProofsToken $proofForToken, ProofsCode $proofForCode, Authorization $authorization) { /** @var Appwrite\Utopia\Database\Documents\User $userFromRequest */ $userFromRequest = $authorization->skip(fn () => $dbForProject->getDocument('users', $userId)); @@ -311,7 +317,7 @@ $createSession = function (string $userId, string $secret, Request $request, Res ]) !== 1; if ($isAllowedTokenType && $hasUserEmail && $isSessionAlertsEnabled && $isNotFirstSession) { - sendSessionAlert($locale, $user, $project, $platform, $session, $queueForMails); + sendSessionAlert($locale, $user, $project, $platform, $session, $publisherForMails); } $queueForEvents @@ -975,13 +981,13 @@ Http::post('/v1/account/sessions/email') ->inject('locale') ->inject('geodb') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('hooks') ->inject('store') ->inject('proofForPassword') ->inject('proofForToken') ->inject('authorization') - ->action(function (string $email, string $password, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, Reader $geodb, Event $queueForEvents, Mail $queueForMails, Hooks $hooks, Store $store, ProofsPassword $proofForPassword, ProofsToken $proofForToken, Authorization $authorization) { + ->action(function (string $email, string $password, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, Reader $geodb, Event $queueForEvents, MailsPublisher $publisherForMails, Hooks $hooks, Store $store, ProofsPassword $proofForPassword, ProofsToken $proofForToken, Authorization $authorization) { $email = \strtolower($email); $protocol = $request->getProtocol(); @@ -1086,7 +1092,7 @@ Http::post('/v1/account/sessions/email') Query::equal('userId', [$user->getId()]), ]) !== 1 ) { - sendSessionAlert($locale, $user, $project, $platform, $session, $queueForMails); + sendSessionAlert($locale, $user, $project, $platform, $session, $publisherForMails); } } @@ -1280,7 +1286,7 @@ Http::post('/v1/account/sessions/token') ->inject('locale') ->inject('geodb') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('store') ->inject('proofForToken') ->inject('proofForCode') @@ -2121,11 +2127,11 @@ Http::post('/v1/account/tokens/magic-url') ->inject('dbForProject') ->inject('locale') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('proofForPassword') ->inject('platform') ->inject('authorization') - ->action(function (string $userId, string $email, string $url, bool $phrase, Request $request, Response $response, Document $user, Document $project, Database $dbForProject, Locale $locale, Event $queueForEvents, Mail $queueForMails, ProofsPassword $proofForPassword, array $platform, Authorization $authorization) { + ->action(function (string $userId, string $email, string $url, bool $phrase, Request $request, Response $response, Document $user, Document $project, Database $dbForProject, Locale $locale, Event $queueForEvents, MailsPublisher $publisherForMails, ProofsPassword $proofForPassword, array $platform, Authorization $authorization) { if (empty(System::getEnv('_APP_SMTP_HOST'))) { throw new Exception(Exception::GENERAL_SMTP_DISABLED, 'SMTP disabled'); } @@ -2281,6 +2287,9 @@ Http::post('/v1/account/tokens/magic-url') $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -2292,12 +2301,13 @@ Http::post('/v1/account/tokens/magic-url') $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -2314,10 +2324,9 @@ Http::post('/v1/account/tokens/magic-url') $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $projectName = $project->getAttribute('name'); @@ -2339,18 +2348,21 @@ Http::post('/v1/account/tokens/magic-url') 'team' => '', ]; - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->appendVariables($emailVariables) - ->setRecipient($email); - if ($project->getId() === 'console') { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $email, + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + customMailOptions: $customMailOptions, + platform: $platform, + )); $token->setAttribute('secret', $tokenSecret); @@ -2401,11 +2413,11 @@ Http::post('/v1/account/tokens/email') ->inject('dbForProject') ->inject('locale') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('proofForPassword') ->inject('proofForCode') ->inject('authorization') - ->action(function (string $userId, string $email, bool $phrase, Request $request, Response $response, User $user, Document $project, array $platform, Database $dbForProject, Locale $locale, Event $queueForEvents, Mail $queueForMails, ProofsPassword $proofForPassword, ProofsCode $proofForCode, Authorization $authorization) { + ->action(function (string $userId, string $email, bool $phrase, Request $request, Response $response, User $user, Document $project, array $platform, Database $dbForProject, Locale $locale, Event $queueForEvents, MailsPublisher $publisherForMails, ProofsPassword $proofForPassword, ProofsCode $proofForCode, Authorization $authorization) { if (empty(System::getEnv('_APP_SMTP_HOST'))) { throw new Exception(Exception::GENERAL_SMTP_DISABLED, 'SMTP disabled'); } @@ -2567,6 +2579,9 @@ Http::post('/v1/account/tokens/email') $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -2578,12 +2593,13 @@ Http::post('/v1/account/tokens/email') $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -2600,10 +2616,9 @@ Http::post('/v1/account/tokens/email') $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $projectName = $project->getAttribute('name'); @@ -2639,20 +2654,22 @@ Http::post('/v1/account/tokens/email') ]); } - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->setBodyTemplate($bodyTemplate) - ->appendVariables($emailVariables) - ->setRecipient($email); - - // since this is console project, set email sender name! if ($smtpBaseTemplate === APP_BRANDED_EMAIL_BASE_TEMPLATE) { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $email, + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + bodyTemplate: $bodyTemplate, + customMailOptions: $customMailOptions, + platform: $platform, + )); $token->setAttribute('secret', $tokenSecret); @@ -2708,14 +2725,14 @@ Http::put('/v1/account/sessions/magic-url') ->inject('locale') ->inject('geodb') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('store') ->inject('proofForCode') ->inject('authorization') - ->action(function ($userId, $secret, $request, $response, $user, $dbForProject, $project, $platform, $locale, $geodb, $queueForEvents, $queueForMails, $store, $proofForCode, $authorization) use ($createSession) { + ->action(function ($userId, $secret, $request, $response, $user, $dbForProject, $project, $platform, $locale, $geodb, $queueForEvents, $publisherForMails, $store, $proofForCode, $authorization) use ($createSession) { $proofForToken = new ProofsToken(TOKEN_LENGTH_MAGIC_URL); $proofForToken->setHash(new Sha()); - $createSession($userId, $secret, $request, $response, $user, $dbForProject, $project, $platform, $locale, $geodb, $queueForEvents, $queueForMails, $store, $proofForToken, $proofForCode, $authorization); + $createSession($userId, $secret, $request, $response, $user, $dbForProject, $project, $platform, $locale, $geodb, $queueForEvents, $publisherForMails, $store, $proofForToken, $proofForCode, $authorization); }); Http::put('/v1/account/sessions/phone') @@ -2757,7 +2774,7 @@ Http::put('/v1/account/sessions/phone') ->inject('locale') ->inject('geodb') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('store') ->inject('proofForToken') ->inject('proofForCode') @@ -3579,11 +3596,11 @@ Http::post('/v1/account/recovery') ->inject('project') ->inject('platform') ->inject('locale') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('queueForEvents') ->inject('proofForToken') ->inject('authorization') - ->action(function (string $email, string $url, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, Mail $queueForMails, Event $queueForEvents, ProofsToken $proofForToken, Authorization $authorization) { + ->action(function (string $email, string $url, Request $request, Response $response, User $user, Database $dbForProject, Document $project, array $platform, Locale $locale, MailsPublisher $publisherForMails, Event $queueForEvents, ProofsToken $proofForToken, Authorization $authorization) { if (empty(System::getEnv('_APP_SMTP_HOST'))) { throw new Exception(Exception::GENERAL_SMTP_DISABLED, 'SMTP Disabled'); @@ -3664,6 +3681,9 @@ Http::post('/v1/account/recovery') $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -3675,12 +3695,13 @@ Http::post('/v1/account/recovery') $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -3697,10 +3718,9 @@ Http::post('/v1/account/recovery') $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $emailVariables = [ @@ -3713,19 +3733,22 @@ Http::post('/v1/account/recovery') 'team' => '' ]; - $queueForMails - ->setRecipient($profile->getAttribute('email', '')) - ->setName($profile->getAttribute('name', '')) - ->setBody($body) - ->appendVariables($emailVariables) - ->setSubject($subject) - ->setPreview($preview); - if ($project->getId() === 'console') { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $profile->getAttribute('email', ''), + name: $profile->getAttribute('name', ''), + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + customMailOptions: $customMailOptions, + platform: $platform, + )); $recovery->setAttribute('secret', $secret); @@ -3893,10 +3916,10 @@ Http::post('/v1/account/verifications/email') ->inject('dbForProject') ->inject('locale') ->inject('queueForEvents') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('proofForToken') ->inject('authorization') - ->action(function (string $url, Request $request, Response $response, Document $project, array $platform, User $user, Database $dbForProject, Locale $locale, Event $queueForEvents, Mail $queueForMails, ProofsToken $proofForToken, Authorization $authorization) { + ->action(function (string $url, Request $request, Response $response, Document $project, array $platform, User $user, Database $dbForProject, Locale $locale, Event $queueForEvents, MailsPublisher $publisherForMails, ProofsToken $proofForToken, Authorization $authorization) { if (empty(System::getEnv('_APP_SMTP_HOST'))) { throw new Exception(Exception::GENERAL_SMTP_DISABLED, 'SMTP Disabled'); @@ -3981,6 +4004,9 @@ Http::post('/v1/account/verifications/email') $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -3992,12 +4018,13 @@ Http::post('/v1/account/verifications/email') $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -4014,10 +4041,9 @@ Http::post('/v1/account/verifications/email') $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $emailVariables = [ @@ -4044,20 +4070,23 @@ Http::post('/v1/account/verifications/email') ]); } - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->setBodyTemplate($bodyTemplate) - ->appendVariables($emailVariables) - ->setRecipient($user->getAttribute('email')) - ->setName($user->getAttribute('name') ?? ''); - if ($project->getId() === 'console') { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $user->getAttribute('email'), + name: $user->getAttribute('name') ?? '', + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + bodyTemplate: $bodyTemplate, + customMailOptions: $customMailOptions, + platform: $platform, + )); $verification->setAttribute('secret', $verificationSecret); diff --git a/app/controllers/api/projects.php b/app/controllers/api/projects.php index 24a1b28cdd..208edf0fcc 100644 --- a/app/controllers/api/projects.php +++ b/app/controllers/api/projects.php @@ -3,7 +3,8 @@ use Ahc\Jwt\JWT; use Appwrite\Auth\Validator\MockNumber; use Appwrite\Event\Delete; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Validator\Event; use Appwrite\Extend\Exception; use Appwrite\Network\Platform; @@ -1837,9 +1838,9 @@ Http::post('/v1/projects/:projectId/smtp/tests') ->param('secure', '', new WhiteList(['tls', 'ssl'], true), 'Does SMTP server use secure connection', true) ->inject('response') ->inject('dbForPlatform') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('plan') - ->action(function (string $projectId, array $emails, string $senderName, string $senderEmail, string $replyTo, string $host, int $port, string $username, string $password, string $secure, Response $response, Database $dbForPlatform, Mail $queueForMails, array $plan) { + ->action(function (string $projectId, array $emails, string $senderName, string $senderEmail, string $replyTo, string $host, int $port, string $username, string $password, string $secure, Response $response, Database $dbForPlatform, MailsPublisher $publisherForMails, array $plan) { $project = $dbForPlatform->getDocument('projects', $projectId); if ($project->isEmpty()) { @@ -1862,22 +1863,25 @@ Http::post('/v1/projects/:projectId/smtp/tests') ->setParam('{{privacyUrl}}', $plan['privacyUrl'] ?? APP_EMAIL_PRIVACY_URL); foreach ($emails as $email) { - $queueForMails - ->setSmtpHost($host) - ->setSmtpPort($port) - ->setSmtpUsername($username) - ->setSmtpPassword($password) - ->setSmtpSecure($secure) - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName) - ->setRecipient($email) - ->setName('') - ->setBodyTemplate(__DIR__ . '/../../config/locale/templates/email-base-styled.tpl') - ->setBody($template->render()) - ->setVariables([]) - ->setSubject($subject) - ->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $email, + name: '', + subject: $subject, + body: $template->render(), + smtp: [ + 'host' => $host, + 'port' => $port, + 'username' => $username, + 'password' => $password, + 'secure' => $secure, + 'replyTo' => $replyTo, + 'senderEmail' => $senderEmail, + 'senderName' => $senderName, + ], + variables: [], + bodyTemplate: __DIR__ . '/../../config/locale/templates/email-base-styled.tpl', + )); } return $response->noContent(); diff --git a/app/controllers/shared/api.php b/app/controllers/shared/api.php index c8824d3708..7583818b06 100644 --- a/app/controllers/shared/api.php +++ b/app/controllers/shared/api.php @@ -9,9 +9,9 @@ use Appwrite\Event\Database as EventDatabase; use Appwrite\Event\Delete; use Appwrite\Event\Event; use Appwrite\Event\Func; -use Appwrite\Event\Mail; use Appwrite\Event\Message\Usage as UsageMessage; use Appwrite\Event\Messaging; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Event\Realtime; use Appwrite\Event\Webhook; @@ -456,7 +456,7 @@ Http::init() ->inject('queueForBuilds') ->inject('usage') ->inject('queueForFunctions') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('dbForProject') ->inject('timelimit') ->inject('resourceToken') @@ -467,7 +467,7 @@ Http::init() ->inject('telemetry') ->inject('platform') ->inject('authorization') - ->action(function (Http $utopia, Request $request, Response $response, Document $project, Document $user, Event $queueForEvents, Messaging $queueForMessaging, Audit $queueForAudits, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Context $usage, Func $queueForFunctions, Mail $queueForMails, Database $dbForProject, callable $timelimit, Document $resourceToken, string $mode, ?Key $apiKey, array $plan, Document $devKey, Telemetry $telemetry, array $platform, Authorization $authorization) { + ->action(function (Http $utopia, Request $request, Response $response, Document $project, Document $user, Event $queueForEvents, Messaging $queueForMessaging, Audit $queueForAudits, Delete $queueForDeletes, EventDatabase $queueForDatabase, Build $queueForBuilds, Context $usage, Func $queueForFunctions, MailsPublisher $publisherForMails, Database $dbForProject, callable $timelimit, Document $resourceToken, string $mode, ?Key $apiKey, array $plan, Document $devKey, Telemetry $telemetry, array $platform, Authorization $authorization) { $route = $utopia->getRoute(); @@ -577,12 +577,10 @@ Http::init() $queueForMessaging->setProject($project); $queueForFunctions->setProject($project); $queueForBuilds->setProject($project); - $queueForMails->setProject($project); /* Auto-set platforms */ $queueForFunctions->setPlatform($platform); $queueForBuilds->setPlatform($platform); - $queueForMails->setPlatform($platform); $useCache = $route->getLabel('cache', false); $storageCacheOperationsCounter = $telemetry->createCounter('storage.cache.operations.load'); diff --git a/app/init/resources.php b/app/init/resources.php index 1bab4491a4..b041073a81 100644 --- a/app/init/resources.php +++ b/app/init/resources.php @@ -11,9 +11,9 @@ use Appwrite\Event\Database as EventDatabase; use Appwrite\Event\Delete; use Appwrite\Event\Event; use Appwrite\Event\Func; -use Appwrite\Event\Mail; use Appwrite\Event\Messaging; use Appwrite\Event\Migration; +use Appwrite\Event\Publisher\Mail; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Event\Realtime; use Appwrite\Event\Screenshot; @@ -126,9 +126,10 @@ Http::setResource('publisherWebhooks', function (Publisher $publisher) { Http::setResource('queueForMessaging', function (Publisher $publisher) { return new Messaging($publisher); }, ['publisher']); -Http::setResource('queueForMails', function (Publisher $publisher) { - return new Mail($publisher); -}, ['publisher']); +Http::setResource('publisherForMails', fn (Publisher $publisher) => new Mail( + $publisher, + new Queue(System::getEnv('_APP_MAILS_QUEUE_NAME', Event::MAILS_QUEUE_NAME)) +), ['publisher']); Http::setResource('queueForBuilds', function (Publisher $publisher) { return new Build($publisher); }, ['publisher']); diff --git a/app/worker.php b/app/worker.php index db036b6a99..79081a391a 100644 --- a/app/worker.php +++ b/app/worker.php @@ -10,9 +10,9 @@ use Appwrite\Event\Database as EventDatabase; use Appwrite\Event\Delete; use Appwrite\Event\Event; use Appwrite\Event\Func; -use Appwrite\Event\Mail; use Appwrite\Event\Messaging; use Appwrite\Event\Migration; +use Appwrite\Event\Publisher\Mail; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Event\Realtime; use Appwrite\Event\Screenshot; @@ -325,9 +325,10 @@ Server::setResource('queueForMessaging', function (Publisher $publisher) { return new Messaging($publisher); }, ['publisher']); -Server::setResource('queueForMails', function (Publisher $publisher) { - return new Mail($publisher); -}, ['publisher']); +Server::setResource('publisherForMails', fn (Publisher $publisher) => new Mail( + $publisher, + new Queue(System::getEnv('_APP_MAILS_QUEUE_NAME', Event::MAILS_QUEUE_NAME)) +), ['publisher']); Server::setResource('queueForBuilds', function (Publisher $publisher) { return new Build($publisher); diff --git a/src/Appwrite/Event/Message/Mail.php b/src/Appwrite/Event/Message/Mail.php new file mode 100644 index 0000000000..99fa0df2fd --- /dev/null +++ b/src/Appwrite/Event/Message/Mail.php @@ -0,0 +1,61 @@ + $this->project, + 'recipient' => $this->recipient, + 'name' => $this->name, + 'subject' => $this->subject, + 'bodyTemplate' => $this->bodyTemplate, + 'body' => $this->body, + 'preview' => $this->preview, + 'smtp' => $this->smtp, + 'variables' => $this->variables, + 'attachment' => $this->attachment, + 'customMailOptions' => $this->customMailOptions, + 'platform' => empty($this->platform) ? Config::getParam('platform', []) : $this->platform, + ]; + } + + public static function fromArray(array $data): static + { + return new static( + project: !empty($data['project']) ? new Document($data['project']) : null, + recipient: $data['recipient'] ?? '', + name: $data['name'] ?? '', + subject: $data['subject'] ?? '', + body: $data['body'] ?? '', + preview: $data['preview'] ?? '', + smtp: $data['smtp'] ?? [], + variables: $data['variables'] ?? [], + bodyTemplate: $data['bodyTemplate'] ?? '', + attachment: $data['attachment'] ?? [], + customMailOptions: $data['customMailOptions'] ?? [], + platform: $data['platform'] ?? [], + ); + } +} diff --git a/src/Appwrite/Event/Publisher/Mail.php b/src/Appwrite/Event/Publisher/Mail.php new file mode 100644 index 0000000000..66c6af7bf5 --- /dev/null +++ b/src/Appwrite/Event/Publisher/Mail.php @@ -0,0 +1,27 @@ +publish($this->queue, $message); + } + + public function getSize(bool $failed = false): int + { + return $this->getQueueSize($this->queue, $failed); + } +} diff --git a/src/Appwrite/Platform/Modules/Account/Http/Account/MFA/Challenges/Create.php b/src/Appwrite/Platform/Modules/Account/Http/Account/MFA/Challenges/Create.php index 20a6afed2e..b064228462 100644 --- a/src/Appwrite/Platform/Modules/Account/Http/Account/MFA/Challenges/Create.php +++ b/src/Appwrite/Platform/Modules/Account/Http/Account/MFA/Challenges/Create.php @@ -5,8 +5,9 @@ namespace Appwrite\Platform\Modules\Account\Http\Account\MFA\Challenges; use Appwrite\Auth\MFA\Type; use Appwrite\Detector\Detector; use Appwrite\Event\Event; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; use Appwrite\Event\Messaging; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Extend\Exception; use Appwrite\SDK\AuthType; use Appwrite\SDK\ContentType; @@ -102,7 +103,7 @@ class Create extends Action ->inject('request') ->inject('queueForEvents') ->inject('queueForMessaging') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('timelimit') ->inject('usage') ->inject('plan') @@ -122,7 +123,7 @@ class Create extends Action Request $request, Event $queueForEvents, Messaging $queueForMessaging, - Mail $queueForMails, + MailsPublisher $publisherForMails, callable $timelimit, Context $usage, array $plan, @@ -255,6 +256,9 @@ class Create extends Action $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ""; + $smtpConfig = []; + $customMailOptions = []; + if ($smtpEnabled) { if (!empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -266,12 +270,13 @@ class Create extends Action $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (!empty($customTemplate)) { if (!empty($customTemplate['senderEmail'])) { @@ -288,10 +293,9 @@ class Create extends Action $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $emailVariables = [ @@ -318,20 +322,22 @@ class Create extends Action ]); } - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->setBodyTemplate($bodyTemplate) - ->appendVariables($emailVariables) - ->setRecipient($user->getAttribute('email')); - - // since this is console project, set email sender name! if ($smtpBaseTemplate === APP_BRANDED_EMAIL_BASE_TEMPLATE) { - $queueForMails->setSenderName($platform['emailSenderName']); + $customMailOptions['senderName'] = $platform['emailSenderName']; } - $queueForMails->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $user->getAttribute('email'), + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + bodyTemplate: $bodyTemplate, + customMailOptions: $customMailOptions, + platform: $platform, + )); break; } diff --git a/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Failed/Get.php b/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Failed/Get.php index cb3640746f..072adae919 100644 --- a/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Failed/Get.php +++ b/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Failed/Get.php @@ -9,9 +9,9 @@ use Appwrite\Event\Database; use Appwrite\Event\Delete; use Appwrite\Event\Event; use Appwrite\Event\Func; -use Appwrite\Event\Mail; use Appwrite\Event\Messaging; use Appwrite\Event\Migration; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Event\Screenshot; use Appwrite\Event\StatsResources; @@ -76,7 +76,7 @@ class Get extends Base ->inject('queueForDatabase') ->inject('queueForDeletes') ->inject('queueForAudits') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('queueForFunctions') ->inject('queueForStatsResources') ->inject('publisherForUsage') @@ -96,7 +96,7 @@ class Get extends Base Database $queueForDatabase, Delete $queueForDeletes, Audit $queueForAudits, - Mail $queueForMails, + MailsPublisher $publisherForMails, Func $queueForFunctions, StatsResources $queueForStatsResources, UsagePublisher $publisherForUsage, @@ -113,7 +113,7 @@ class Get extends Base System::getEnv('_APP_DATABASE_QUEUE_NAME', Event::DATABASE_QUEUE_NAME) => $queueForDatabase, System::getEnv('_APP_DELETE_QUEUE_NAME', Event::DELETE_QUEUE_NAME) => $queueForDeletes, System::getEnv('_APP_AUDITS_QUEUE_NAME', Event::AUDITS_QUEUE_NAME) => $queueForAudits, - System::getEnv('_APP_MAILS_QUEUE_NAME', Event::MAILS_QUEUE_NAME) => $queueForMails, + System::getEnv('_APP_MAILS_QUEUE_NAME', Event::MAILS_QUEUE_NAME) => $publisherForMails, System::getEnv('_APP_FUNCTIONS_QUEUE_NAME', Event::FUNCTIONS_QUEUE_NAME) => $queueForFunctions, System::getEnv('_APP_STATS_RESOURCES_QUEUE_NAME', Event::STATS_RESOURCES_QUEUE_NAME) => $queueForStatsResources, System::getEnv('_APP_STATS_USAGE_QUEUE_NAME', Event::STATS_USAGE_QUEUE_NAME) => $publisherForUsage, diff --git a/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Mails/Get.php b/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Mails/Get.php index 3b9c06b5f9..cf8a470c1a 100644 --- a/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Mails/Get.php +++ b/src/Appwrite/Platform/Modules/Health/Http/Health/Queue/Mails/Get.php @@ -2,7 +2,7 @@ namespace Appwrite\Platform\Modules\Health\Http\Health\Queue\Mails; -use Appwrite\Event\Mail; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Platform\Modules\Health\Http\Health\Queue\Base; use Appwrite\SDK\AuthType; use Appwrite\SDK\ContentType; @@ -42,16 +42,16 @@ class Get extends Base contentType: ContentType::JSON )) ->param('threshold', 5000, new Integer(true), 'Queue size threshold. When hit (equal or higher), endpoint returns server error. Default value is 5000.', true) - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('response') ->callback($this->action(...)); } - public function action(int|string $threshold, Mail $queueForMails, Response $response): void + public function action(int|string $threshold, MailsPublisher $publisherForMails, Response $response): void { $threshold = (int) $threshold; - $size = $queueForMails->getSize(); + $size = $publisherForMails->getSize(); $this->assertQueueThreshold($size, $threshold); diff --git a/src/Appwrite/Platform/Modules/Teams/Http/Memberships/Create.php b/src/Appwrite/Platform/Modules/Teams/Http/Memberships/Create.php index 221c3aa521..1243bd6e5c 100644 --- a/src/Appwrite/Platform/Modules/Teams/Http/Memberships/Create.php +++ b/src/Appwrite/Platform/Modules/Teams/Http/Memberships/Create.php @@ -4,8 +4,9 @@ namespace Appwrite\Platform\Modules\Teams\Http\Memberships; use Appwrite\Auth\Validator\Phone; use Appwrite\Event\Event; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; use Appwrite\Event\Messaging; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Extend\Exception; use Appwrite\Network\Validator\Email as EmailValidator; use Appwrite\Platform\Action; @@ -87,7 +88,7 @@ class Create extends Action ->inject('dbForProject') ->inject('authorization') ->inject('locale') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('queueForMessaging') ->inject('queueForEvents') ->inject('timelimit') @@ -98,7 +99,7 @@ class Create extends Action ->callback($this->action(...)); } - public function action(string $teamId, string $email, string $userId, string $phone, array $roles, string $url, string $name, Response $response, Document $project, Document $user, Database $dbForProject, Authorization $authorization, Locale $locale, Mail $queueForMails, Messaging $queueForMessaging, Event $queueForEvents, callable $timelimit, Context $usage, array $plan, Password $proofForPassword, Token $proofForToken) + public function action(string $teamId, string $email, string $userId, string $phone, array $roles, string $url, string $name, Response $response, Document $project, Document $user, Database $dbForProject, Authorization $authorization, Locale $locale, MailsPublisher $publisherForMails, Messaging $queueForMessaging, Event $queueForEvents, callable $timelimit, Context $usage, array $plan, Password $proofForPassword, Token $proofForToken) { $isAppUser = User::isApp($authorization->getRoles()); $isPrivilegedUser = User::isPrivileged($authorization->getRoles()); @@ -317,6 +318,8 @@ class Create extends Action $senderName = System::getEnv('_APP_SYSTEM_EMAIL_NAME', APP_NAME . ' Server'); $replyTo = ''; + $smtpConfig = []; + if ($smtpEnabled) { if (! empty($smtp['senderEmail'])) { $senderEmail = $smtp['senderEmail']; @@ -328,12 +331,13 @@ class Create extends Action $replyTo = $smtp['replyTo']; } - $queueForMails - ->setSmtpHost($smtp['host'] ?? '') - ->setSmtpPort($smtp['port'] ?? '') - ->setSmtpUsername($smtp['username'] ?? '') - ->setSmtpPassword($smtp['password'] ?? '') - ->setSmtpSecure($smtp['secure'] ?? ''); + $smtpConfig = [ + 'host' => $smtp['host'] ?? '', + 'port' => $smtp['port'] ?? '', + 'username' => $smtp['username'] ?? '', + 'password' => $smtp['password'] ?? '', + 'secure' => $smtp['secure'] ?? '', + ]; if (! empty($customTemplate)) { if (! empty($customTemplate['senderEmail'])) { @@ -350,10 +354,9 @@ class Create extends Action $subject = $customTemplate['subject'] ?? $subject; } - $queueForMails - ->setSmtpReplyTo($replyTo) - ->setSmtpSenderEmail($senderEmail) - ->setSmtpSenderName($senderName); + $smtpConfig['replyTo'] = $replyTo; + $smtpConfig['senderEmail'] = $senderEmail; + $smtpConfig['senderName'] = $senderName; } $emailVariables = [ @@ -366,14 +369,16 @@ class Create extends Action 'project' => $projectName, ]; - $queueForMails - ->setSubject($subject) - ->setBody($body) - ->setPreview($preview) - ->setRecipient($invitee->getAttribute('email')) - ->setName($invitee->getAttribute('name', '')) - ->appendVariables($emailVariables) - ->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $invitee->getAttribute('email'), + name: $invitee->getAttribute('name', ''), + subject: $subject, + body: $body, + preview: $preview, + smtp: $smtpConfig, + variables: $emailVariables, + )); } elseif (! empty($phone)) { if (empty(System::getEnv('_APP_SMS_PROVIDER'))) { throw new Exception(Exception::GENERAL_PHONE_DISABLED, 'Phone provider not configured'); diff --git a/src/Appwrite/Platform/Workers/Certificates.php b/src/Appwrite/Platform/Workers/Certificates.php index 73509819a9..fef9826158 100644 --- a/src/Appwrite/Platform/Workers/Certificates.php +++ b/src/Appwrite/Platform/Workers/Certificates.php @@ -6,7 +6,8 @@ use Appwrite\Certificates\Adapter as CertificatesAdapter; use Appwrite\Event\Certificate; use Appwrite\Event\Event; use Appwrite\Event\Func; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Realtime; use Appwrite\Event\Webhook; use Appwrite\Extend\Exception as AppwriteException; @@ -29,7 +30,7 @@ use Utopia\Database\Validator\Authorization as ValidatorAuthorization; use Utopia\Domains\Domain; use Utopia\Locale\Locale; use Utopia\Logger\Log; -use Utopia\Queue\Message; +use Utopia\Queue\Message as QueueMessage; use Utopia\System\System; class Certificates extends Action @@ -50,7 +51,7 @@ class Certificates extends Action ->desc('Certificates worker') ->inject('message') ->inject('dbForPlatform') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('queueForEvents') ->inject('queueForWebhooks') ->inject('queueForFunctions') @@ -64,9 +65,9 @@ class Certificates extends Action } /** - * @param Message $message + * @param QueueMessage $message * @param Database $dbForPlatform - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param Event $queueForEvents * @param Webhook $queueForWebhooks * @param Func $queueForFunctions @@ -81,9 +82,9 @@ class Certificates extends Action * @throws \Utopia\Database\Exception */ public function action( - Message $message, + QueueMessage $message, Database $dbForPlatform, - Mail $queueForMails, + MailsPublisher $publisherForMails, Event $queueForEvents, Webhook $queueForWebhooks, Func $queueForFunctions, @@ -115,7 +116,7 @@ class Certificates extends Action break; case Certificate::ACTION_GENERATION: - $this->handleCertificateGenerationAction($domain, $domainType, $dbForPlatform, $queueForMails, $queueForEvents, $queueForWebhooks, $queueForFunctions, $queueForRealtime, $log, $certificates, $authorization, $skipRenewCheck, $plan, $validationDomain); + $this->handleCertificateGenerationAction($domain, $domainType, $dbForPlatform, $publisherForMails, $queueForEvents, $queueForWebhooks, $queueForFunctions, $queueForRealtime, $log, $certificates, $authorization, $skipRenewCheck, $plan, $validationDomain); break; default: @@ -204,7 +205,7 @@ class Certificates extends Action * @param Domain $domain * @param ?string $domainType * @param Database $dbForPlatform - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param Event $queueForEvents * @param Webhook $queueForWebhooks * @param Func $queueForFunctions @@ -228,7 +229,7 @@ class Certificates extends Action Domain $domain, ?string $domainType, Database $dbForPlatform, - Mail $queueForMails, + MailsPublisher $publisherForMails, Event $queueForEvents, Webhook $queueForWebhooks, Func $queueForFunctions, @@ -353,7 +354,7 @@ class Certificates extends Action $rule->setAttribute('status', RULE_STATUS_CERTIFICATE_GENERATION_FAILED); // Send email to security email - $this->notifyError($domain->get(), $e->getMessage(), $attempts, $queueForMails, $plan); + $this->notifyError($domain->get(), $e->getMessage(), $attempts, $publisherForMails, $plan); throw $e; } finally { @@ -519,12 +520,12 @@ class Certificates extends Action * @param string $domain Domain that caused the error * @param string $errorMessage Verbose error message * @param int $attempt How many times it failed already - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param array $plan * @return void * @throws Exception */ - private function notifyError(string $domain, string $errorMessage, int $attempt, Mail $queueForMails, array $plan): void + private function notifyError(string $domain, string $errorMessage, int $attempt, MailsPublisher $publisherForMails, array $plan): void { // Log error into console Console::warning('Cannot renew domain (' . $domain . ') on attempt no. ' . $attempt . ' certificate: ' . $errorMessage); @@ -555,14 +556,14 @@ class Certificates extends Action $subject = $locale->getText("emails.certificate.subject"); $preview = $locale->getText("emails.certificate.preview"); - $queueForMails - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body) - ->setName('Appwrite Administrator') - ->setBodyTemplate(__DIR__ . '/../../../../app/config/locale/templates/email-base-styled.tpl') - ->setVariables($emailVariables) - ->setRecipient(System::getEnv('_APP_EMAIL_CERTIFICATES', System::getEnv('_APP_SYSTEM_SECURITY_EMAIL_ADDRESS'))) - ->trigger(); + $publisherForMails->enqueue(new Mail( + recipient: System::getEnv('_APP_EMAIL_CERTIFICATES', System::getEnv('_APP_SYSTEM_SECURITY_EMAIL_ADDRESS')), + name: 'Appwrite Administrator', + subject: $subject, + body: $body, + preview: $preview, + variables: $emailVariables, + bodyTemplate: __DIR__ . '/../../../../app/config/locale/templates/email-base-styled.tpl', + )); } } diff --git a/src/Appwrite/Platform/Workers/Migrations.php b/src/Appwrite/Platform/Workers/Migrations.php index ce20358626..860075e080 100644 --- a/src/Appwrite/Platform/Workers/Migrations.php +++ b/src/Appwrite/Platform/Workers/Migrations.php @@ -3,8 +3,9 @@ namespace Appwrite\Platform\Workers; use Ahc\Jwt\JWT; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; use Appwrite\Event\Message\Usage as UsageMessage; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Event\Realtime; use Appwrite\Extend\Exception; @@ -41,7 +42,7 @@ use Utopia\Migration\Sources\NHost; use Utopia\Migration\Sources\Supabase; use Utopia\Migration\Transfer; use Utopia\Platform\Action; -use Utopia\Queue\Message; +use Utopia\Queue\Message as QueueMessage; use Utopia\Storage\Device; use Utopia\System\System; @@ -85,7 +86,7 @@ class Migrations extends Action ->inject('queueForRealtime') ->inject('deviceForMigrations') ->inject('deviceForFiles') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('usage') ->inject('publisherForUsage') ->inject('plan') @@ -97,7 +98,7 @@ class Migrations extends Action * @throws Exception */ public function action( - Message $message, + QueueMessage $message, Document $project, Database $dbForProject, Database $dbForPlatform, @@ -105,7 +106,7 @@ class Migrations extends Action Realtime $queueForRealtime, Device $deviceForMigrations, Device $deviceForFiles, - Mail $queueForMails, + MailsPublisher $publisherForMails, Context $usage, UsagePublisher $publisherForUsage, array $plan, @@ -150,7 +151,7 @@ class Migrations extends Action $this->processMigration( $migration, $queueForRealtime, - $queueForMails, + $publisherForMails, $usage, $publisherForUsage, $platform, @@ -349,7 +350,7 @@ class Migrations extends Action protected function processMigration( Document $migration, Realtime $queueForRealtime, - Mail $queueForMails, + MailsPublisher $publisherForMails, Context $usage, UsagePublisher $publisherForUsage, array $platform, @@ -519,7 +520,7 @@ class Migrations extends Action // TODO: Move to CSV hook if ($migration->getAttribute('destination') === DestinationCSV::getName()) { - $this->handleCSVExportComplete($project, $migration, $queueForMails, $queueForRealtime, $platform, $authorization); + $this->handleCSVExportComplete($project, $migration, $publisherForMails, $queueForRealtime, $platform, $authorization); } } } finally { @@ -538,7 +539,7 @@ class Migrations extends Action * * @param Document $project * @param Document $migration - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param Realtime $queueForRealtime * @param array $platform * @param Authorization $authorization @@ -547,7 +548,7 @@ class Migrations extends Action protected function handleCSVExportComplete( Document $project, Document $migration, - Mail $queueForMails, + MailsPublisher $publisherForMails, Realtime $queueForRealtime, array $platform, Authorization $authorization, @@ -598,7 +599,7 @@ class Migrations extends Action project: $project, user: $user, options: $options, - queueForMails: $queueForMails, + publisherForMails: $publisherForMails, platform: $platform, sizeMB: $sizeMB ); @@ -660,7 +661,7 @@ class Migrations extends Action project: $project, user: $user, options: $options, - queueForMails: $queueForMails, + publisherForMails: $publisherForMails, platform: $platform, downloadUrl: $downloadUrl ); @@ -673,7 +674,7 @@ class Migrations extends Action * @param Document $project * @param Document $user The user who triggered the operation * @param array $options Migration options - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param array $platform * @param string $downloadUrl Download URL for successful exports * @param float $sizeMB File size in MB for failed exports @@ -685,7 +686,7 @@ class Migrations extends Action Document $project, Document $user, array $options, - Mail $queueForMails, + MailsPublisher $publisherForMails, array $platform, string $downloadUrl = '', float $sizeMB = 0.0, @@ -752,17 +753,20 @@ class Migrations extends Action 'platform' => $platform['platformName'], ]; - $queueForMails - ->setProject($project) - ->setSubject($subject) - ->setPreview($preview) - ->setBody($emailBody) - ->setBodyTemplate(__DIR__ . '/../../../../app/config/locale/templates/email-base-styled.tpl') - ->setVariables($emailVariables) - ->setName($user->getAttribute('name', $user->getAttribute('email'))) - ->setRecipient($user->getAttribute('email')) - ->setSenderName($platform['emailSenderName']) - ->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $user->getAttribute('email'), + name: $user->getAttribute('name', $user->getAttribute('email')), + subject: $subject, + body: $emailBody, + preview: $preview, + variables: $emailVariables, + bodyTemplate: __DIR__ . '/../../../../app/config/locale/templates/email-base-styled.tpl', + customMailOptions: [ + 'senderName' => $platform['emailSenderName'], + ], + platform: $platform, + )); Console::info("CSV export {$emailType} notification email sent to " . $user->getAttribute('email')); } diff --git a/src/Appwrite/Platform/Workers/Webhooks.php b/src/Appwrite/Platform/Workers/Webhooks.php index fce3c7b149..2c9f150498 100644 --- a/src/Appwrite/Platform/Workers/Webhooks.php +++ b/src/Appwrite/Platform/Workers/Webhooks.php @@ -2,8 +2,9 @@ namespace Appwrite\Platform\Workers; -use Appwrite\Event\Mail; +use Appwrite\Event\Message\Mail; use Appwrite\Event\Message\Usage as UsageMessage; +use Appwrite\Event\Publisher\Mail as MailsPublisher; use Appwrite\Event\Publisher\Usage as UsagePublisher; use Appwrite\Template\Template; use Appwrite\Usage\Context as UsageContext; @@ -13,7 +14,7 @@ use Utopia\Database\Document; use Utopia\Database\Query; use Utopia\Logger\Log; use Utopia\Platform\Action; -use Utopia\Queue\Message; +use Utopia\Queue\Message as QueueMessage; use Utopia\System\System; class Webhooks extends Action @@ -36,7 +37,7 @@ class Webhooks extends Action ->inject('message') ->inject('project') ->inject('dbForPlatform') - ->inject('queueForMails') + ->inject('publisherForMails') ->inject('publisherForUsage') ->inject('log') ->inject('plan') @@ -44,17 +45,17 @@ class Webhooks extends Action } /** - * @param Message $message + * @param QueueMessage $message * @param Document $project * @param Database $dbForPlatform - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param UsagePublisher $publisherForUsage * @param Log $log * @param array $plan * @return void * @throws Exception */ - public function action(Message $message, Document $project, Database $dbForPlatform, Mail $queueForMails, UsagePublisher $publisherForUsage, Log $log, array $plan): void + public function action(QueueMessage $message, Document $project, Database $dbForPlatform, MailsPublisher $publisherForMails, UsagePublisher $publisherForUsage, Log $log, array $plan): void { $this->errors = []; $payload = $message->getPayload() ?? []; @@ -73,7 +74,7 @@ class Webhooks extends Action foreach ($project->getAttribute('webhooks', []) as $webhook) { if (array_intersect($webhook->getAttribute('events', []), $events)) { - $this->execute($events, $webhookPayload, $webhook, $user, $project, $dbForPlatform, $queueForMails, $publisherForUsage, $plan); + $this->execute($events, $webhookPayload, $webhook, $user, $project, $dbForPlatform, $publisherForMails, $publisherForUsage, $plan); } } @@ -89,11 +90,11 @@ class Webhooks extends Action * @param Document $user * @param Document $project * @param Database $dbForPlatform - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param array $plan * @return void */ - private function execute(array $events, string $payload, Document $webhook, Document $user, Document $project, Database $dbForPlatform, Mail $queueForMails, UsagePublisher $publisherForUsage, array $plan): void + private function execute(array $events, string $payload, Document $webhook, Document $user, Document $project, Database $dbForPlatform, MailsPublisher $publisherForMails, UsagePublisher $publisherForUsage, array $plan): void { if ($webhook->getAttribute('enabled') !== true) { return; @@ -175,7 +176,7 @@ class Webhooks extends Action if ($attempts >= \intval(System::getEnv('_APP_WEBHOOK_MAX_FAILED_ATTEMPTS', '10'))) { $webhook->setAttribute('enabled', false); $updatePayload['enabled'] = false; - $this->sendEmailAlert($attempts, $statusCode, $webhook, $project, $dbForPlatform, $queueForMails, $plan); + $this->sendEmailAlert($attempts, $statusCode, $webhook, $project, $dbForPlatform, $publisherForMails, $plan); } $dbForPlatform->updateDocument('webhooks', $webhook->getId(), new Document($updatePayload)); @@ -207,11 +208,11 @@ class Webhooks extends Action * @param Document $webhook * @param Document $project * @param Database $dbForPlatform - * @param Mail $queueForMails + * @param MailsPublisher $publisherForMails * @param array $plan * @return void */ - public function sendEmailAlert(int $attempts, mixed $statusCode, Document $webhook, Document $project, Database $dbForPlatform, Mail $queueForMails, array $plan): void + public function sendEmailAlert(int $attempts, mixed $statusCode, Document $webhook, Document $project, Database $dbForPlatform, MailsPublisher $publisherForMails, array $plan): void { $memberships = $dbForPlatform->find('memberships', [ Query::equal('teamInternalId', [$project->getAttribute('teamInternalId')]), @@ -255,18 +256,16 @@ class Webhooks extends Action ->setParam('{{message}}', $template->render()) ->setParam('{{year}}', date("Y")); - $queueForMails - ->setProject($project) - ->setSubject($subject) - ->setPreview($preview) - ->setBody($body->render()); - foreach ($users as $user) { - $queueForMails - ->setVariables(['user' => $user->getAttribute('name', '')]) - ->setName($user->getAttribute('name', '')) - ->setRecipient($user->getAttribute('email')) - ->trigger(); + $publisherForMails->enqueue(new Mail( + project: $project, + recipient: $user->getAttribute('email'), + name: $user->getAttribute('name', ''), + subject: $subject, + body: $body->render(), + preview: $preview, + variables: ['user' => $user->getAttribute('name', '')], + )); } } }