Refactor mails queue to stateless publisher

This commit is contained in:
loks0n
2026-03-15 19:00:32 +00:00
parent 5a56a2fc35
commit b9ce7373a4
14 changed files with 413 additions and 277 deletions
+151 -122
View File
@@ -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);
+23 -19
View File
@@ -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();
+3 -5
View File
@@ -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');
+5 -4
View File
@@ -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']);
+5 -4
View File
@@ -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);
+61
View File
@@ -0,0 +1,61 @@
<?php
namespace Appwrite\Event\Message;
use Utopia\Config\Config;
use Utopia\Database\Document;
readonly class Mail extends Base
{
public function __construct(
public ?Document $project = null,
public string $recipient = '',
public string $name = '',
public string $subject = '',
public string $body = '',
public string $preview = '',
public array $smtp = [],
public array $variables = [],
public string $bodyTemplate = '',
public array $attachment = [],
public array $customMailOptions = [],
public array $platform = [],
) {
}
public function toArray(): array
{
return [
'project' => $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'] ?? [],
);
}
}
+27
View File
@@ -0,0 +1,27 @@
<?php
namespace Appwrite\Event\Publisher;
use Appwrite\Event\Message\Mail as Message;
use Utopia\Queue\Publisher;
use Utopia\Queue\Queue;
readonly class Mail extends Base
{
public function __construct(
Publisher $publisher,
protected Queue $queue
) {
parent::__construct($publisher);
}
public function enqueue(Message $message): string|bool
{
return $this->publish($this->queue, $message);
}
public function getSize(bool $failed = false): int
{
return $this->getQueueSize($this->queue, $failed);
}
}
@@ -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;
}
@@ -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,
@@ -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);
@@ -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');
+23 -22
View File
@@ -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',
));
}
}
+29 -25
View File
@@ -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'));
}
+22 -23
View File
@@ -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', '')],
));
}
}
}