This commit is contained in:
Jake Barnby
2025-04-17 17:09:08 +12:00
parent 100870557b
commit 70364d7a07
9 changed files with 31 additions and 32 deletions
+3 -2
View File
@@ -2,8 +2,8 @@
namespace Appwrite\Messaging\Adapter;
use Appwrite\PubSub\Adapter as PubSubAdapter;
use Appwrite\Messaging\Adapter as MessagingAdapter;
use Appwrite\PubSub\Adapter as PubSubAdapter;
use Utopia\Database\DateTime;
use Utopia\Database\Document;
use Utopia\Database\Helpers\ID;
@@ -142,7 +142,8 @@ class Realtime extends MessagingAdapter
global $register;
$register->get('pools')->get('pubsub')->use(fn (PubSubAdapter $redis) =>
$register->get('pools')->get('pubsub')->use(
fn (PubSubAdapter $redis) =>
$redis->publish('realtime', json_encode([
'project' => $projectId,
'roles' => $roles,
@@ -60,7 +60,7 @@ class ScheduleExecutions extends ScheduleBase
\go(function () use ($schedule, $delay, $data, $pools) {
Co::sleep($delay);
$pools->get('publisher')->use(function(Publisher $publisher) use ($schedule, $data) {
$pools->get('publisher')->use(function (Publisher $publisher) use ($schedule, $data) {
$queueForFunctions = new Func($publisher);
$queueForFunctions->setType('schedule')
@@ -84,6 +84,6 @@ class ScheduleExecutions extends ScheduleBase
);
unset($this->schedules[$schedule['$internalId']]);
}
}
}
}
@@ -75,28 +75,28 @@ class ScheduleFunctions extends ScheduleBase
\sleep($delay); // in seconds
foreach ($scheduleKeys as $scheduleKey) {
// Ensure schedule was not deleted
if (!\array_key_exists($scheduleKey, $this->schedules)) {
return;
}
$schedule = $this->schedules[$scheduleKey];
$this->updateProjectAccess($schedule['project'], $dbForPlatform);
$pools->get('publisher')->use(function(Publisher $publisher) use ($schedule) {
$queueForFunctions = new Func($publisher);
$queueForFunctions
->setType('schedule')
->setFunction($schedule['resource'])
->setMethod('POST')
->setPath('/')
->setProject($schedule['project'])
->trigger();
});
foreach ($scheduleKeys as $scheduleKey) {
// Ensure schedule was not deleted
if (!\array_key_exists($scheduleKey, $this->schedules)) {
return;
}
$schedule = $this->schedules[$scheduleKey];
$this->updateProjectAccess($schedule['project'], $dbForPlatform);
$pools->get('publisher')->use(function (Publisher $publisher) use ($schedule) {
$queueForFunctions = new Func($publisher);
$queueForFunctions
->setType('schedule')
->setFunction($schedule['resource'])
->setMethod('POST')
->setPath('/')
->setProject($schedule['project'])
->trigger();
});
}
});
}
@@ -42,7 +42,7 @@ class ScheduleMessages extends ScheduleBase
}
\go(function () use ($schedule, $pools, $dbForPlatform) {
$pools->get('publisher')->use(function(Publisher $publisher) use ($schedule, $dbForPlatform) {
$pools->get('publisher')->use(function (Publisher $publisher) use ($schedule, $dbForPlatform) {
$queueForMessaging = new Messaging($publisher);
$this->updateProjectAccess($schedule['project'], $dbForPlatform);