mirror of
https://github.com/appwrite/appwrite.git
synced 2026-05-26 13:51:13 +00:00
updated the queue for event propagation
This commit is contained in:
+13
-13
@@ -322,14 +322,8 @@ if (!function_exists('triggerPresenceUsage')) {
|
||||
if (!function_exists('getQueueForEventsForProject')) {
|
||||
function getQueueForEventsForProject(Document $project, User $user): QueueEvent
|
||||
{
|
||||
global $register;
|
||||
|
||||
/** @var Group $pools */
|
||||
$pools = $register->get('pools');
|
||||
|
||||
$queueForEvents = new QueueEvent(new BrokerPool(
|
||||
publisher: $pools->get('publisher')
|
||||
));
|
||||
global $container;
|
||||
$queueForEvents = $container->get('queueForEvents');
|
||||
|
||||
$queueForEvents->setProject($project);
|
||||
$queueForEvents->setUser($user);
|
||||
@@ -340,8 +334,6 @@ if (!function_exists('getQueueForEventsForProject')) {
|
||||
|
||||
if (!function_exists('triggerPresenceEvent')) {
|
||||
function triggerPresenceEvent(
|
||||
Server $server,
|
||||
Realtime $realtime,
|
||||
Document $project,
|
||||
User $user,
|
||||
string $eventName,
|
||||
@@ -407,9 +399,17 @@ if (!function_exists('setPermission')) {
|
||||
}
|
||||
|
||||
global $container;
|
||||
|
||||
$container->set('pools', function ($register) {
|
||||
return $register->get('pools');
|
||||
}, ['register']);
|
||||
|
||||
$container->set('queueForEvents', function ($pools) {
|
||||
return new QueueEvent(new BrokerPool(
|
||||
publisher: $pools->get('publisher')
|
||||
));
|
||||
}, ['pools']);
|
||||
|
||||
$container->set('queueForRealtime', function () {
|
||||
return new QueueRealtime();
|
||||
}, []);
|
||||
@@ -1559,7 +1559,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re
|
||||
}
|
||||
}
|
||||
|
||||
triggerPresenceEvent($server, $realtime, $project, $user, 'presences.[presenceId].upsert', $presence);
|
||||
triggerPresenceEvent($project, $user, 'presences.[presenceId].upsert', $presence);
|
||||
|
||||
break;
|
||||
|
||||
@@ -1613,7 +1613,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re
|
||||
}
|
||||
});
|
||||
|
||||
$server->onClose(function (int $connection) use ($server, $realtime, $stats, $register) {
|
||||
$server->onClose(function (int $connection) use ($realtime, $stats, $register) {
|
||||
$projectId = null;
|
||||
$userId = null;
|
||||
$subscriptionsBeforeClose = 0;
|
||||
@@ -1670,7 +1670,7 @@ $server->onClose(function (int $connection) use ($server, $realtime, $stats, $re
|
||||
|
||||
foreach ($presences as $presence) {
|
||||
try {
|
||||
triggerPresenceEvent($server, $realtime, $project, new User([]), 'presences.[presenceId].delete', $presence);
|
||||
triggerPresenceEvent($project, new User([]), 'presences.[presenceId].delete', $presence);
|
||||
} catch (Throwable) {
|
||||
// Swallow errors to avoid breaking disconnect cleanup
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user