From 0e3a017d5008b28cf61c229bf84f021db8ec16f2 Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Thu, 16 Apr 2026 12:35:50 +0530 Subject: [PATCH] added realtime trigger --- app/realtime.php | 75 +++++++++++++++++-- .../Presence/PresenceRealtimeClientTest.php | 48 ++++++++++++ 2 files changed, 118 insertions(+), 5 deletions(-) diff --git a/app/realtime.php b/app/realtime.php index adf5294e9b..acd15e7540 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -2,6 +2,8 @@ use Appwrite\Extend\Exception; use Appwrite\Extend\Exception as AppwriteException; +use Appwrite\Event\Event as QueueEvent; +use Appwrite\Event\Realtime as QueueRealtime; use Appwrite\Messaging\Adapter\Realtime; use Appwrite\Network\Validator\Origin; use Appwrite\PubSub\Adapter\Pool as PubSubPool; @@ -38,6 +40,7 @@ use Utopia\DI\Container; use Utopia\DSN\DSN; use Utopia\Logger\Log; use Utopia\Pools\Group; +use Utopia\Queue\Broker\Pool as BrokerPool; use Utopia\Registry\Registry; use Utopia\System\System; use Utopia\Telemetry\Adapter\None as NoTelemetry; @@ -224,6 +227,7 @@ if (!function_exists('getRealtime')) { } } + if (!function_exists('getTelemetry')) { function getTelemetry(int $workerId): Utopia\Telemetry\Adapter { @@ -243,6 +247,57 @@ if (!function_exists('triggerStats')) { } } +if (!function_exists('triggerPresenceEvent')) { + 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') + )); + + $queueForEvents->setProject($project); + $queueForEvents->setUser($user); + + return $queueForEvents; + } + + function triggerPresenceEvent( + Server $server, + Realtime $realtime, + 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()); + + (new QueueRealtime()) + ->setProject($project) + ->setUser($user) + ->from($queueForEvents) + ->trigger(); + } catch (Throwable $th) { + logError($th, 'realtimePresenceEvent', tags: [ + 'projectId' => $project->getId(), + 'event' => $eventName, + ]); + } + } +} + if (!function_exists('setPermission')) { function setPermission(Document $document, ?array $permissions, Authorization $authorization): void { @@ -1181,6 +1236,8 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re } } + triggerPresenceEvent($server, $realtime, $project, $user, 'presences.[presenceId].upsert', $presence); + break; default: @@ -1216,7 +1273,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re } }); -$server->onClose(function (int $connection) use ($realtime, $stats, $register) { +$server->onClose(function (int $connection) use ($server, $realtime, $stats, $register) { try { if (array_key_exists($connection, $realtime->connections)) { $stats->decr($realtime->connections[$connection]['projectId'], 'connectionsTotal'); @@ -1236,13 +1293,21 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) { $consoleDB = getConsoleDB(); $project = $consoleDB->getAuthorization()->skip(fn () => $consoleDB->getDocument('projects', $projectId)); + // todo: have a bulk delete if (!$project->isEmpty()) { $dbForProject = getProjectDB($project); - $dbForProject->deleteDocuments('presenceLogs', [ + $presences = $dbForProject->find('presenceLogs', [ Query::equal('$id', $presenceIds), - ], onError: function (Throwable $th) { - // Swallow errors to avoid breaking disconnect cleanup - }); + ]); + + foreach ($presences as $presence) { + try { + $dbForProject->deleteDocument('presenceLogs', $presence->getId()); + triggerPresenceEvent($server, $realtime, $project, new User([]), 'presences.[presenceId].delete', $presence); + } catch (Throwable) { + // Swallow errors to avoid breaking disconnect cleanup + } + } } } } diff --git a/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php b/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php index 299a12cadb..9889e84105 100644 --- a/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php +++ b/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php @@ -285,4 +285,52 @@ class PresenceRealtimeClientTest extends Scope $client->close(); } + + public function testPresenceMessageEmitsCreateAndDeleteEvents(): void + { + $presenceId = ID::unique(); + $userId = $this->getUser()['$id']; + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $this->getProject()['$id'] . '=' . $this->getUser()['session'], + ]; + + $listener = $this->getWebsocket(['presences', 'presences.' . $presenceId], $headers, timeout: 8); + $connected = \json_decode($listener->receive(), true); + $this->assertSame('connected', $connected['type'] ?? null); + + $publisher = $this->connectPresenceSocket(true, timeout: 8); + + $publisher->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $presenceId, + 'status' => 'online', + 'metadata' => [ + 'source' => 'realtime-create-delete-events', + ], + 'permissions' => $this->getPresencePermissions($userId), + ], + ])); + + $createResponse = \json_decode($publisher->receive(), true); + $this->assertSame('response', $createResponse['type'] ?? null); + $this->assertSame('presence', $createResponse['data']['to'] ?? null); + $this->assertSame($presenceId, $createResponse['data']['presence']['$id'] ?? null); + + $createEvent = \json_decode($listener->receive(), true); + $this->assertSame('event', $createEvent['type'] ?? null); + $this->assertContains('presences.' . $presenceId . '.upsert', $createEvent['data']['events'] ?? []); + $this->assertSame($presenceId, $createEvent['data']['payload']['$id'] ?? null); + $this->assertSame('online', $createEvent['data']['payload']['status'] ?? null); + + $publisher->close(); + + $deleteEvent = \json_decode($listener->receive(), true); + $this->assertSame('event', $deleteEvent['type'] ?? null); + $this->assertContains('presences.' . $presenceId . '.delete', $deleteEvent['data']['events'] ?? []); + $this->assertSame($presenceId, $deleteEvent['data']['payload']['$id'] ?? null); + + $listener->close(); + } }