diff --git a/app/realtime.php b/app/realtime.php index 98e90e594d..6361cd7a0b 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -1,7 +1,6 @@ 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 } diff --git a/src/Appwrite/Presences/State.php b/src/Appwrite/Presences/State.php index becded6074..19e7dc98b7 100644 --- a/src/Appwrite/Presences/State.php +++ b/src/Appwrite/Presences/State.php @@ -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, + ]); + } + } } } diff --git a/src/Appwrite/Realtime/Message/Handlers/Presence.php b/src/Appwrite/Realtime/Message/Handlers/Presence.php index d3bf65b0e3..1787bc0682 100644 --- a/src/Appwrite/Realtime/Message/Handlers/Presence.php +++ b/src/Appwrite/Realtime/Message/Handlers/Presence.php @@ -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|null $permissions - * @param Closure(int, Document): void $triggerPresenceUsage - * @param Closure(?Document, User, string, Document): void $triggerPresenceEvent * @return array */ 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',