added realtime trigger

This commit is contained in:
ArnabChatterjee20k
2026-04-16 12:35:50 +05:30
parent d28cce761d
commit 0e3a017d50
2 changed files with 118 additions and 5 deletions
+70 -5
View File
@@ -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
}
}
}
}
}
@@ -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();
}
}