feat: enhance presence handling by refactoring usage and event triggering logic

This commit is contained in:
ArnabChatterjee20k
2026-05-14 11:48:41 +05:30
parent 0aa8d402ae
commit a3542ed7fa
3 changed files with 117 additions and 100 deletions
+24 -75
View File
@@ -1,7 +1,6 @@
<?php
use Appwrite\Event\Event as QueueEvent;
use Appwrite\Event\Message\Usage as UsageMessage;
use Appwrite\Event\Publisher\Usage as UsagePublisher;
use Appwrite\Event\Realtime as QueueRealtime;
use Appwrite\Extend\Exception;
@@ -16,7 +15,6 @@ use Appwrite\Realtime\Message\Handlers\Ping as PingHandler;
use Appwrite\Realtime\Message\Handlers\Presence as PresenceHandler;
use Appwrite\Realtime\Message\Handlers\Subscribe as SubscribeHandler;
use Appwrite\Realtime\Message\Handlers\Unsubscribe as UnsubscribeHandler;
use Appwrite\Usage\Context;
use Appwrite\Utopia\Database\Documents\User;
use Appwrite\Utopia\Request;
use Appwrite\Utopia\Response;
@@ -259,8 +257,8 @@ if (!function_exists('getTelemetry')) {
}
}
if (!function_exists('getQueueForEventsForProject')) {
function getQueueForEventsForProject(Document $project, User $user): QueueEvent
if (!function_exists('getQueueForEvents')) {
function getQueueForEvents(): QueueEvent
{
$ctx = Coroutine::getContext();
@@ -273,10 +271,7 @@ if (!function_exists('getQueueForEventsForProject')) {
));
}
return $ctx['queueForEvents']
->reset()
->setProject($project)
->setUser($user);
return $ctx['queueForEvents'];
}
}
@@ -289,7 +284,7 @@ if (!function_exists('getQueueForRealtime')) {
$ctx['queueForRealtime'] = new QueueRealtime();
}
return $ctx['queueForRealtime']->reset();
return $ctx['queueForRealtime'];
}
}
@@ -299,63 +294,6 @@ if (!function_exists('triggerStats')) {
}
}
if (!function_exists('triggerPresenceUsage')) {
function triggerPresenceUsage(int $value, Document $project): void
{
if ($project->isEmpty()) {
return;
}
try {
global $container;
/** @var UsagePublisher $publisherForUsage */
$publisherForUsage = $container->get('publisherForUsage');
$usage = new Context();
$usage->addMetric(METRIC_USERS_PRESENCE, $value);
$publisherForUsage->enqueue(new UsageMessage(
project: $project,
metrics: $usage->getMetrics(),
));
} catch (Throwable $th) {
logError($th, 'realtimeStats', tags: ['projectId' => $project->getId()]);
}
}
}
if (!function_exists('triggerPresenceEvent')) {
function triggerPresenceEvent(
Document $project,
User $user,
string $eventName,
Document $presence
): void {
if ($project->isEmpty() || $presence->isEmpty()) {
return;
}
try {
$queueForEvents = getQueueForEventsForProject($project, $user);
$queueForEvents
->setEvent($eventName)
->setParam('presenceId', $presence->getId())
->setPayload($presence->getArrayCopy());
getQueueForRealtime()
->setProject($project)
->setUser($user)
->from($queueForEvents)
->trigger();
} catch (Throwable $th) {
logError($th, 'realtimePresenceEvent', tags: [
'projectId' => $project->getId(),
'event' => $eventName,
]);
}
}
}
global $container;
if (!$container->has('pools')) {
@@ -1197,11 +1135,9 @@ $server->onMessage(function (int $connection, string $message) use ($container,
$messageContainer->set('authorization', fn () => $authorization);
$messageContainer->set('project', fn () => $project);
$messageContainer->set('projectId', fn () => $projectId);
// Wrap the global helpers as first-class callables so handlers receive them via
// injection rather than reaching into the global namespace. Keeps PresenceHandler
// unit-testable without bootstrapping app/realtime.php.
$messageContainer->set('triggerPresenceUsage', fn () => triggerPresenceUsage(...));
$messageContainer->set('triggerPresenceEvent', fn () => triggerPresenceEvent(...));
$messageContainer->set('publisherForUsage', fn () => $container->get('publisherForUsage'));
$messageContainer->set('queueForEvents', fn () => getQueueForEvents());
$messageContainer->set('queueForRealtime', fn () => getQueueForRealtime());
$responsePayload = $messageDispatcher->dispatch($messageContainer, $message);
@@ -1265,7 +1201,7 @@ $server->onMessage(function (int $connection, string $message) use ($container,
}
});
$server->onClose(function (int $connection) use ($realtime, $stats, $register) {
$server->onClose(function (int $connection) use ($realtime, $stats, $register, $container, $presenceState) {
$projectId = null;
$userId = null;
$subscriptionsBeforeClose = 0;
@@ -1297,7 +1233,7 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
!empty($presencesById)
&& $projectId !== 'console'
) {
go(function () use ($presencesById, $projectId, $userId): void {
go(function () use ($presencesById, $projectId, $userId, $container, $presenceState): void {
// Fresh span: the parent realtime.close span finishes before this coroutine
Span::init('realtime.close.presenceCleanup');
Span::add('realtime.projectId', $projectId);
@@ -1329,9 +1265,12 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
}
}
/** @var UsagePublisher $publisherForUsage */
$publisherForUsage = $container->get('publisherForUsage');
try {
$deletionCount = $dbForProject->deleteDocuments('presenceLogs', [Query::equal('$id', $presenceIds)]);
triggerPresenceUsage(-$deletionCount, $project);
$presenceState->triggerUsage($publisherForUsage, $project, -$deletionCount);
} catch (Throwable $th) {
Span::error($th);
logError($th, 'realtimeOnClosePresenceDeletion', tags: [
@@ -1340,9 +1279,19 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
]);
}
$queueForEvents = getQueueForEvents();
$queueForRealtime = getQueueForRealtime();
foreach ($presences as $presence) {
try {
triggerPresenceEvent($project, $user, 'presences.[presenceId].delete', $presence);
$presenceState->triggerEvent(
$queueForEvents,
$queueForRealtime,
$project,
$user,
'presences.[presenceId].delete',
$presence,
);
} catch (Throwable) {
// Swallow errors to avoid breaking disconnect cleanup
}
+74 -15
View File
@@ -2,8 +2,14 @@
namespace Appwrite\Presences;
use Appwrite\Event\Event as QueueEvent;
use Appwrite\Event\Message\Usage as UsageMessage;
use Appwrite\Event\Publisher\Usage as UsagePublisher;
use Appwrite\Event\Realtime as QueueRealtime;
use Appwrite\Extend\Exception;
use Appwrite\Usage\Context as UsageContext;
use Appwrite\Utopia\Database\Documents\User;
use Throwable;
use Utopia\Database\Database;
use Utopia\Database\Document;
use Utopia\Database\Exception\Conflict as ConflictException;
@@ -150,18 +156,6 @@ class State
}
}
private function getListCacheKey(Database $dbForProject): string
{
return \sprintf(
'%s-cache:%s:%s:%s:collection:%s',
$dbForProject->getCacheName(),
$dbForProject->getAdapter()->getHostname(),
$dbForProject->getNamespace(),
$dbForProject->getTenant(),
self::COLLECTION_ID
);
}
private function getListCacheFieldKey(array $roles, array $queries, string $type): string
{
$serialized = \array_map(
@@ -185,9 +179,10 @@ class State
int $ttl
): mixed {
$cacheField = $this->getListCacheFieldKey($roles, $queries, $type);
[$collectionKey] = $dbForProject->getCacheKeys(self::COLLECTION_ID);
try {
return $dbForProject->getCache()->load($this->getListCacheKey($dbForProject), $ttl, $cacheField);
return $dbForProject->getCache()->load($collectionKey, $ttl, $cacheField);
} catch (\Throwable) {
return null;
}
@@ -201,15 +196,79 @@ class State
mixed $value
): void {
$cacheField = $this->getListCacheFieldKey($roles, $queries, $type);
[$collectionKey] = $dbForProject->getCacheKeys(self::COLLECTION_ID);
try {
$dbForProject->getCache()->save($this->getListCacheKey($dbForProject), $value, $cacheField);
$dbForProject->getCache()->save($collectionKey, $value, $cacheField);
} catch (\Throwable) {
}
}
public function purgeListCache(Database $dbForProject): bool
{
return $dbForProject->getCache()->purge($this->getListCacheKey($dbForProject));
[$collectionKey] = $dbForProject->getCacheKeys(self::COLLECTION_ID);
return $dbForProject->getCache()->purge($collectionKey);
}
public function triggerUsage(
UsagePublisher $publisher,
Document $project,
int $value,
): void {
if ($project->isEmpty()) {
return;
}
try {
$usage = new UsageContext();
$usage->addMetric(METRIC_USERS_PRESENCE, $value);
$publisher->enqueue(new UsageMessage(
project: $project,
metrics: $usage->getMetrics(),
));
} catch (Throwable $th) {
if (\function_exists('logError')) {
\logError($th, 'realtimeStats', tags: ['projectId' => $project->getId()]);
}
}
}
public function triggerEvent(
QueueEvent $queueForEvents,
QueueRealtime $queueForRealtime,
Document $project,
User $user,
string $eventName,
Document $presence,
): void {
if ($project->isEmpty() || $presence->isEmpty()) {
return;
}
try {
$queueForEvents
->reset()
->setProject($project)
->setUser($user)
->setEvent($eventName)
->setParam('presenceId', $presence->getId())
->setPayload($presence->getArrayCopy());
$queueForRealtime
->reset()
->setProject($project)
->setUser($user)
->from($queueForEvents)
->trigger();
} catch (Throwable $th) {
if (\function_exists('logError')) {
\logError($th, 'realtimePresenceEvent', tags: [
'projectId' => $project->getId(),
'event' => $eventName,
]);
}
}
}
}
@@ -2,12 +2,14 @@
namespace Appwrite\Realtime\Message\Handlers;
use Appwrite\Event\Event as QueueEvent;
use Appwrite\Event\Publisher\Usage as UsagePublisher;
use Appwrite\Event\Realtime as QueueRealtime;
use Appwrite\Extend\Exception;
use Appwrite\Messaging\Adapter\Realtime;
use Appwrite\Presences\State as PresenceState;
use Appwrite\Realtime\Message\Dispatcher;
use Appwrite\Utopia\Database\Documents\User;
use Closure;
use Utopia\Database\Database;
use Utopia\Database\DateTime;
use Utopia\Database\Document;
@@ -34,15 +36,14 @@ class Presence extends Action
->inject('authorization')
->inject('presenceState')
->inject('project')
->inject('triggerPresenceUsage')
->inject('triggerPresenceEvent')
->inject('publisherForUsage')
->inject('queueForEvents')
->inject('queueForRealtime')
->callback($this->action(...));
}
/**
* @param array<int, string>|null $permissions
* @param Closure(int, Document): void $triggerPresenceUsage
* @param Closure(?Document, User, string, Document): void $triggerPresenceEvent
* @return array<string, mixed>
*/
public function action(
@@ -56,8 +57,9 @@ class Presence extends Action
Authorization $authorization,
PresenceState $presenceState,
?Document $project,
Closure $triggerPresenceUsage,
Closure $triggerPresenceEvent,
UsagePublisher $publisherForUsage,
QueueEvent $queueForEvents,
QueueRealtime $queueForRealtime,
): array {
if ($project === null || $project->isEmpty()) {
throw new Exception(Exception::REALTIME_POLICY_VIOLATION, 'Presence requires a project context.');
@@ -93,8 +95,8 @@ class Presence extends Action
$presenceDocument,
$presenceId,
(string) $user->getSequence(),
function () use ($project, $triggerPresenceUsage): void {
$triggerPresenceUsage(1, $project);
function () use ($presenceState, $publisherForUsage, $project): void {
$presenceState->triggerUsage($publisherForUsage, $project, 1);
},
);
@@ -106,7 +108,14 @@ class Presence extends Action
$realtime->connections[$connectionId]['presences'][$presence->getId()] = $presence;
$triggerPresenceEvent($project, $user, 'presences.[presenceId].upsert', $presence);
$presenceState->triggerEvent(
$queueForEvents,
$queueForRealtime,
$project,
$user,
'presences.[presenceId].upsert',
$presence,
);
return [
'type' => 'response',