feat: update presence handling and refactor related attributes for consistency

This commit is contained in:
ArnabChatterjee20k
2026-05-13 18:52:14 +05:30
parent f818f16ebe
commit f09aec7651
9 changed files with 109 additions and 99 deletions
+6 -6
View File
@@ -2828,7 +2828,7 @@ return [
],
[
'$id' => ID::custom('metadata'),
'type' => Database::VAR_STRING,
'type' => Database::VAR_TEXT,
'format' => '',
'size' => 65535,
'signed' => true,
@@ -2838,7 +2838,7 @@ return [
'filters' => ['json'],
],
[
'$id' => ID::custom('perms_md5'),
'$id' => ID::custom('permissionsHash'),
'type' => Database::VAR_STRING,
'format' => '',
'size' => 32,
@@ -2888,13 +2888,13 @@ return [
'$id' => ID::custom('_key_source_status'),
'type' => Database::INDEX_KEY,
'attributes' => ['source', 'status'],
'lengths' => [Database::LENGTH_KEY, Database::LENGTH_KEY],
'orders' => [Database::ORDER_ASC, Database::ORDER_ASC]
'lengths' => [Database::LENGTH_KEY],
'orders' => [Database::ORDER_ASC]
],
[
'$id' => ID::custom('_key_perms_md5'),
'$id' => ID::custom('_key_permissionsHash'),
'type' => Database::INDEX_KEY,
'attributes' => ['perms_md5'],
'attributes' => ['permissionsHash'],
'lengths' => [32],
'orders' => [Database::ORDER_ASC]
]
+85 -64
View File
@@ -64,29 +64,6 @@ if (System::getEnv('_APP_EDITION', 'self-hosted') === 'self-hosted') {
/** @var Registry $register */
$register = $GLOBALS['register'] ?? throw new \RuntimeException('Registry not initialized');
/** @var Group $pools */
$pools = $register->get('pools');
$statsUsageConnection = System::getEnv('_APP_CONNECTIONS_QUEUE_STATS_USAGE', '');
$publisherPoolName = 'publisher';
if (!empty($statsUsageConnection)) {
try {
$pools->get('publisher_' . $statsUsageConnection);
$publisherPoolName = 'publisher_' . $statsUsageConnection;
} catch (Throwable) {
// Fallback to default publisher pool when custom one is unavailable.
}
}
$broker = new BrokerPool(publisher: $pools->get($publisherPoolName));
$publisherForUsage = new UsagePublisher(
$broker,
new Queue(System::getEnv(
'_APP_STATS_USAGE_QUEUE_NAME',
\Appwrite\Event\Event::STATS_USAGE_QUEUE_NAME
))
);
$registerConnectionResources ??= require __DIR__ . '/init/realtime/connection.php';
Runtime::enableCoroutine(SWOOLE_HOOK_ALL);
@@ -282,35 +259,6 @@ if (!function_exists('getTelemetry')) {
}
}
if (!function_exists('triggerStats')) {
function triggerStats(array $event, string $projectId): void
{
}
}
if (!function_exists('triggerPresenceUsage')) {
function triggerPresenceUsage(int $value, Document $project): void
{
if ($project->isEmpty()) {
return;
}
try {
global $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('getQueueForEventsForProject')) {
function getQueueForEventsForProject(Document $project, User $user): QueueEvent
{
@@ -345,6 +293,37 @@ if (!function_exists('getQueueForRealtime')) {
}
}
if (!function_exists('triggerStats')) {
function triggerStats(array $event, string $projectId): void
{
}
}
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,
@@ -379,19 +358,45 @@ if (!function_exists('triggerPresenceEvent')) {
global $container;
$container->set('pools', function ($register) {
return $register->get('pools');
}, ['register']);
if (!$container->has('pools')) {
$container->set('pools', function ($register) {
return $register->get('pools');
}, ['register']);
}
if (!$container->has('publisherForUsage')) {
$container->set('publisherForUsage', function (Group $pools): UsagePublisher {
$statsUsageConnection = System::getEnv('_APP_CONNECTIONS_QUEUE_STATS_USAGE', '');
$publisherPoolName = 'publisher';
if (!empty($statsUsageConnection)) {
try {
$pools->get('publisher_' . $statsUsageConnection);
$publisherPoolName = 'publisher_' . $statsUsageConnection;
} catch (Throwable) {
// Fallback to default publisher pool when custom one is unavailable.
}
}
return new UsagePublisher(
new BrokerPool(publisher: $pools->get($publisherPoolName)),
new Queue(System::getEnv(
'_APP_STATS_USAGE_QUEUE_NAME',
QueueEvent::STATS_USAGE_QUEUE_NAME
))
);
}, ['pools']);
}
$realtime = getRealtime();
$presenceState = new PresenceState();
$messageDispatcher = (new MessageDispatcher())
->register(new PingHandler())
->register(new AuthenticationHandler())
->register(new SubscribeHandler())
->register(new UnsubscribeHandler())
->register(new PresenceHandler());
->addHandler(new PingHandler())
->addHandler(new AuthenticationHandler())
->addHandler(new SubscribeHandler())
->addHandler(new UnsubscribeHandler())
->addHandler(new PresenceHandler());
/**
* Table for statistics across all workers.
@@ -1331,15 +1336,15 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
!empty($presencesById)
&& $projectId !== 'console'
) {
go(function () use ($presencesById, $projectId): void {
go(function () use ($presencesById, $projectId, $userId): void {
// Fresh span: the parent realtime.close span finishes before this coroutine
Span::init('realtime.close.presenceCleanup');
Span::add('realtime.projectId', $projectId);
Span::add('realtime.presenceCount', \count($presencesById));
try {
$consoleDB = getConsoleDB();
$project = $consoleDB->getAuthorization()->skip(fn () => $consoleDB->getDocument('projects', $projectId));
$dbForPlatform = getConsoleDB();
$project = $dbForPlatform->getAuthorization()->skip(fn () => $dbForPlatform->getDocument('projects', $projectId));
if ($project->isEmpty()) {
return;
@@ -1349,6 +1354,22 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
$presences = \array_values($presencesById);
$dbForProject = getProjectDB($project);
// Resolve the disconnecting user so event payloads carry proper actor context.
$user = new User([]);
if (!empty($userId)) {
try {
/** @var User $fetched */
$fetched = $dbForProject->getAuthorization()->skip(
fn () => $dbForProject->getDocument('users', $userId)
);
if (!$fetched->isEmpty()) {
$user = $fetched;
}
} catch (Throwable) {
// Fall back to empty User if lookup fails.
}
}
try {
$deletionCount = $dbForProject->deleteDocuments('presenceLogs', [Query::equal('$id', $presenceIds)]);
triggerPresenceUsage(-$deletionCount, $project);
@@ -1362,7 +1383,7 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
foreach ($presences as $presence) {
try {
triggerPresenceEvent($project, new User([]), 'presences.[presenceId].delete', $presence);
triggerPresenceEvent($project, $user, 'presences.[presenceId].delete', $presence);
} catch (Throwable) {
// Swallow errors to avoid breaking disconnect cleanup
}
+9 -9
View File
@@ -45,12 +45,12 @@ class PresenceState
}
if (!$isAPIKey && !$isPrivilegedUser) {
$this->assertPermissionsAgainstAuthorization($permissions, $authorization);
$this->checkPermissions($permissions, $authorization);
}
sort($permissions, SORT_STRING);
$document->setAttribute('$permissions', $permissions);
$document->setAttribute('perms_md5', \md5(\json_encode($permissions)));
$document->setAttribute('permissionsHash', \md5(\json_encode($permissions)));
return $document;
}
@@ -146,12 +146,12 @@ class PresenceState
sort($ownerPermissions, SORT_STRING);
$document->setAttribute('$permissions', $ownerPermissions);
$document->setAttribute('perms_md5', \md5(\json_encode($ownerPermissions)));
$document->setAttribute('permissionsHash', \md5(\json_encode($ownerPermissions)));
return $document;
}
private function assertPermissionsAgainstAuthorization(array $permissions, Authorization $authorization): void
private function checkPermissions(array $permissions, Authorization $authorization): void
{
foreach (Database::PERMISSIONS as $type) {
foreach ($permissions as $permission) {
@@ -185,7 +185,7 @@ class PresenceState
);
}
private function getListCacheField(array $roles, array $queries, string $type): string
private function getListCacheFieldKey(array $roles, array $queries, string $type): string
{
$serialized = \array_map(
static fn ($query) => $query instanceof Query ? $query->toArray() : $query,
@@ -200,14 +200,14 @@ class PresenceState
);
}
public function loadListCacheField(
public function getListCacheField(
Database $dbForProject,
array $roles,
array $queries,
string $type,
int $ttl
): mixed {
$cacheField = $this->getListCacheField($roles, $queries, $type);
$cacheField = $this->getListCacheFieldKey($roles, $queries, $type);
try {
return $dbForProject->getCache()->load($this->getListCacheKey($dbForProject), $ttl, $cacheField);
@@ -216,14 +216,14 @@ class PresenceState
}
}
public function saveListCacheField(
public function setListCacheField(
Database $dbForProject,
array $roles,
array $queries,
string $type,
mixed $value
): void {
$cacheField = $this->getListCacheField($roles, $queries, $type);
$cacheField = $this->getListCacheFieldKey($roles, $queries, $type);
try {
$dbForProject->getCache()->save($this->getListCacheKey($dbForProject), $value, $cacheField);
@@ -106,7 +106,7 @@ class XList extends PlatformAction
$roles = $dbForProject->getAuthorization()->getRoles();
$documentsCacheHit = false;
$cachedDocuments = $presenceState->loadListCacheField(
$cachedDocuments = $presenceState->getListCacheField(
$dbForProject,
$roles,
$queries,
@@ -126,7 +126,7 @@ class XList extends PlatformAction
$documentsArray = \array_map(function ($doc) {
return $doc->getArrayCopy();
}, $documents);
$presenceState->saveListCacheField(
$presenceState->setListCacheField(
$dbForProject,
$roles,
$queries,
@@ -136,7 +136,7 @@ class XList extends PlatformAction
}
if ($includeTotal) {
$cachedTotal = $presenceState->loadListCacheField(
$cachedTotal = $presenceState->getListCacheField(
$dbForProject,
$roles,
$filterQueries,
@@ -147,7 +147,7 @@ class XList extends PlatformAction
$total = (int) $cachedTotal;
} else {
$total = $dbForProject->count('presenceLogs', [...$filterQueries, $expiryFilter], APP_LIMIT_COUNT);
$presenceState->saveListCacheField(
$presenceState->setListCacheField(
$dbForProject,
$roles,
$filterQueries,
+1 -1
View File
@@ -22,7 +22,7 @@ class Dispatcher
*/
private array $handlers = [];
public function register(Action $handler): self
public function addHandler(Action $handler): self
{
$labels = $handler->getLabels();
$type = $labels[self::LABEL_MESSAGE_TYPE]
@@ -101,7 +101,7 @@ class Presence extends Action
$presence->removeAttribute('$collection');
$presence->removeAttribute('$tenant');
$presence->removeAttribute('hostname');
$presence->removeAttribute('perms_md5');
$presence->removeAttribute('permissionsHash');
$presence->removeAttribute('userInternalId');
$realtime->connections[$connectionId]['presences'][$presence->getId()] = $presence;
@@ -29,13 +29,6 @@ class Presence extends Any
'default' => '',
'example' => '5e5ea5c16897e',
])
->addRule('$sequence', [
'type' => self::TYPE_ID,
'description' => 'Presence sequence ID.',
'default' => '',
'example' => '1',
'readOnly' => true,
])
->addRule('$createdAt', [
'type' => self::TYPE_DATETIME,
'description' => 'Presence creation date in ISO 8601 format.',
@@ -90,13 +83,9 @@ class Presence extends Any
$document->removeAttribute('$collection');
$document->removeAttribute('$tenant');
$document->removeAttribute('hostname');
$document->removeAttribute('perms_md5');
$document->removeAttribute('permissionsHash');
$document->removeAttribute('userInternalId');
if (!$document->isEmpty()) {
$document->setAttribute('$sequence', (string) $document->getAttribute('$sequence', ''));
}
foreach ($document->getAttributes() as $attribute) {
if (\is_array($attribute)) {
foreach ($attribute as $subAttribute) {
@@ -257,7 +257,7 @@ class RealtimeConsoleClientTest extends Scope
$this->assertEquals('error', $response['type']);
$this->assertNotEmpty($response['data']);
$this->assertEquals(1003, $response['data']['code']);
$this->assertEquals('Payload is not valid. session is required', $response['data']['message']);
$this->assertEquals('Payload is not valid. Session is required', $response['data']['message']);
$client->send(\json_encode([
'type' => 'unknown',
@@ -239,7 +239,7 @@ class RealtimeCustomClientTest extends Scope
$this->assertEquals('error', $response['type']);
$this->assertNotEmpty($response['data']);
$this->assertEquals(1003, $response['data']['code']);
$this->assertEquals('Payload is not valid. session is required', $response['data']['message']);
$this->assertEquals('Payload is not valid. Session is required', $response['data']['message']);
$client->send(\json_encode([
'type' => 'unknown',