diff --git a/app/init.php b/app/init.php index 348285cfdb..bb2b6f2a6d 100644 --- a/app/init.php +++ b/app/init.php @@ -849,10 +849,12 @@ App::setResource('register', fn() => $register); App::setResource('locale', fn() => new Locale(App::getEnv('_APP_LOCALE', 'en'))); // Queues -App::setResource('messaging', fn() => new Phone()); App::setResource('queue', function (Group $pools) { return $pools->get('queue')->pop()->getResource(); }, ['pools']); +App::setResource('messaging', function (Connection $queue) { + return new Phone($queue); +}, ['queue']); App::setResource('mails', function (Connection $queue) { return new Mail($queue); }, ['queue']); diff --git a/app/workers/messaging.php b/app/workers/messaging.php index 5732c8c00b..1c2b96c10b 100644 --- a/app/workers/messaging.php +++ b/app/workers/messaging.php @@ -1,7 +1,9 @@ getUser(); +$secret = $dsn->getPassword(); -class MessagingV1 extends Worker -{ - protected ?Adapter $sms = null; - protected ?string $from = null; +Server::setResource('sms', function () use ($dsn, $user, $secret) { + return match ($dsn->getHost()) { + 'mock' => new Mock($user, $secret), // used for tests + 'twilio' => new Twilio($user, $secret), + 'text-magic' => new TextMagic($user, $secret), + 'telesign' => new Telesign($user, $secret), + 'msg91' => new Msg91($user, $secret), + 'vonage' => new Vonage($user, $secret), + default => null + }; +}); - public function getName(): string - { - return "mails"; - } +Server::setResource('execute', function () { + return function (string $recipient, string $message, Adapter $sms) { + $from = App::getEnv('_APP_SMS_FROM'); - public function init(): void - { - $dsn = new DSN(App::getEnv('_APP_SMS_PROVIDER')); - $user = $dsn->getUser(); - $secret = $dsn->getPassword(); - - $this->sms = match ($dsn->getHost()) { - 'mock' => new Mock($user, $secret), // used for tests - 'twilio' => new Twilio($user, $secret), - 'text-magic' => new TextMagic($user, $secret), - 'telesign' => new Telesign($user, $secret), - 'msg91' => new Msg91($user, $secret), - 'vonage' => new Vonage($user, $secret), - default => null - }; - - $this->from = App::getEnv('_APP_SMS_FROM'); - } - - public function run(): void - { if (empty(App::getEnv('_APP_SMS_PROVIDER'))) { Console::info('Skipped sms processing. No Phone provider has been set.'); return; } - if (empty($this->from)) { + if (empty($from)) { Console::info('Skipped sms processing. No phone number has been set.'); return; } $message = new SMS( - to: [$this->args['recipient']], - content: $this->args['message'], - from: $this->from, + to: [$recipient], + content: $message, + from: $from, ); try { - $this->sms->send($message); + $sms->send($message); } catch (\Exception $error) { throw new Exception('Error sending message: ' . $error->getMessage(), 500); } - } + }; +}); - public function shutdown(): void - { - } -} +$server->job() + ->inject('message') + ->inject('execute') + ->inject('sms') + ->action(function (Message $message, callable $execute, Adapter $sms) { + $payload = $message->getPayload() ?? []; + + if (empty($payload)) { + throw new Exception('Missing payload'); + } + + if (empty($payload['recipient'])) { + throw new Exception('Missing recipient'); + } + + if (empty($payload['message'])) { + throw new Exception('Missing message'); + } + + $execute($payload['recipient'], $payload['message'], $sms); + }); + +$server->workerStart(); +$server->start(); diff --git a/src/Appwrite/Event/Phone.php b/src/Appwrite/Event/Phone.php index 8baa5120c9..0399e27e45 100644 --- a/src/Appwrite/Event/Phone.php +++ b/src/Appwrite/Event/Phone.php @@ -2,14 +2,15 @@ namespace Appwrite\Event; -use Resque; +use Utopia\Queue\Client; +use Utopia\Queue\Connection; class Phone extends Event { protected string $recipient = ''; protected string $message = ''; - public function __construct() + public function __construct(protected Connection $connection) { parent::__construct(Event::MESSAGING_QUEUE_NAME, Event::MESSAGING_CLASS_NAME); } @@ -68,7 +69,11 @@ class Phone extends Event */ public function trigger(): string|bool { - return Resque::enqueue($this->queue, $this->class, [ + $client = new Client($this->queue, $this->connection); + + $events = $this->getEvent() ? Event::generateEvents($this->getEvent(), $this->getParams()) : null; + + return $client->enqueue([ 'project' => $this->project, 'user' => $this->user, 'payload' => $this->payload,