Add triggerStats function for realtime usage metrics tracking

This commit is contained in:
ArnabChatterjee20k
2026-03-05 18:10:09 +05:30
parent 9fdd7c1c6e
commit b5c2cc9716
+100 -60
View File
@@ -251,6 +251,43 @@ if (!function_exists('getQueueForStatsUsageForProject')) {
}
}
if (!function_exists('triggerStats')) {
/**
* Trigger realtime usage stats with a generic metric map.
*
* @param array<string,int|float> $event Metrics in the form METRIC_CONSTANT => value
*/
function triggerStats(array $event, string $projectId): void
{
if (empty($projectId)) {
return;
}
try {
$consoleDB = getConsoleDB();
/** @var Document $project */
$project = $consoleDB->getAuthorization()->skip(
fn () => $consoleDB->getDocument('projects', $projectId)
);
if ($project->isEmpty()) {
return;
}
$queueForStatsUsage = getQueueForStatsUsageForProject($project);
foreach ($event as $metric => $value) {
$queueForStatsUsage->addMetric($metric, $value);
}
$queueForStatsUsage->trigger();
$queueForStatsUsage->reset();
} catch (Throwable $th) {
logError($th, 'realtimeStats', tags: ['projectId' => $projectId]);
}
}
}
$realtime = getRealtime();
/**
@@ -597,33 +634,15 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats,
$projectId = $event['project'] ?? null;
if (!empty($projectId)) {
try {
$consoleDB = getConsoleDB();
/** @var Document $project */
$project = $consoleDB->getAuthorization()->skip(
fn () => $consoleDB->getDocument('projects', $projectId)
);
$metrics = [
METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT => $total,
];
if (!$project->isEmpty()) {
$queueForStatsUsage = getQueueForStatsUsageForProject($project);
$queueForStatsUsage->addMetric(
METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT,
$total
);
if ($outboundBytes > 0) {
$queueForStatsUsage->addMetric(
METRIC_REALTIME_OUTBOUND,
$outboundBytes
);
}
$queueForStatsUsage->trigger();
}
} catch (Throwable $th) {
logError($th, 'realtimeUsageOutbound', tags: ['projectId' => $projectId]);
if ($outboundBytes > 0) {
$metrics[METRIC_REALTIME_OUTBOUND] = $outboundBytes;
}
triggerStats($metrics, $projectId);
}
}
});
@@ -701,6 +720,19 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server,
throw new Exception(Exception::REALTIME_TOO_MANY_MESSAGES, 'Too many requests');
}
// Record realtime inbound bytes for this project (WS handshake + query params)
try {
$rawSize = $request->getSize();
} catch (Throwable) {
$rawSize = \strlen((string) $request->getURI());
}
if ($rawSize > 0) {
triggerStats([
METRIC_REALTIME_INBOUND => $rawSize,
], $project->getId());
}
/*
* Validate Client Domain - Check to avoid CSRF attack.
* Adding Appwrite API domains to allow XDOMAIN communication.
@@ -755,14 +787,16 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server,
$user = empty($user->getId()) ? null : $response->output($user, Response::MODEL_ACCOUNT);
$server->send([$connection], json_encode([
$connectedPayloadJson = json_encode([
'type' => 'connected',
'data' => [
'channels' => $names,
'subscriptions' => $mapping,
'user' => $user
]
]));
]);
$server->send([$connection], $connectedPayloadJson);
$register->get('telemetry.connectionCounter')->add(1);
$register->get('telemetry.connectionCreatedCounter')->add(1);
@@ -774,14 +808,10 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server,
$stats->incr($project->getId(), 'connections');
$stats->incr($project->getId(), 'connectionsTotal');
try {
$queueForStatsUsage = getQueueForStatsUsageForProject($project);
$queueForStatsUsage
->addMetric(METRIC_REALTIME_CONNECTIONS, 1)
->trigger();
} catch (\Throwable $th) {
logError($th, 'realtimeUsageConnections', project: $project);
}
$connectedOutboundBytes = \strlen($connectedPayloadJson);
triggerStats([METRIC_REALTIME_CONNECTIONS => 1, METRIC_REALTIME_OUTBOUND => $connectedOutboundBytes], $project->getId());
} catch (Throwable $th) {
logError($th, 'realtime', project: $project, user: $logUser, authorization: $authorization);
@@ -865,14 +895,9 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re
// Record realtime inbound bytes for this project
if ($project !== null && !$project->isEmpty()) {
try {
$queueForStatsUsage = getQueueForStatsUsageForProject($project);
$queueForStatsUsage
->addMetric(METRIC_REALTIME_INBOUND, $rawSize)
->trigger();
} catch (Throwable $th) {
logError($th, 'realtimeUsageInbound', project: $project);
}
triggerStats([
METRIC_REALTIME_INBOUND => $rawSize,
], $project->getId());
}
$message = json_decode($message, true);
@@ -883,9 +908,21 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re
switch ($message['type']) {
case 'ping':
$server->send([$connection], json_encode([
$pongPayloadJson = json_encode([
'type' => 'pong'
]));
]);
$server->send([$connection], $pongPayloadJson);
if ($project !== null && !$project->isEmpty()) {
$pongOutboundBytes = \strlen($pongPayloadJson);
if ($pongOutboundBytes > 0) {
triggerStats([
METRIC_REALTIME_OUTBOUND => $pongOutboundBytes,
], $project->getId());
}
}
break;
case 'authentication':
@@ -946,14 +983,27 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re
}
$user = $response->output($user, Response::MODEL_ACCOUNT);
$server->send([$connection], json_encode([
$authResponsePayloadJson = json_encode([
'type' => 'response',
'data' => [
'to' => 'authentication',
'success' => true,
'user' => $user
]
]));
]);
$server->send([$connection], $authResponsePayloadJson);
if ($project !== null && !$project->isEmpty()) {
$authOutboundBytes = \strlen($authResponsePayloadJson);
if ($authOutboundBytes > 0) {
triggerStats([
METRIC_REALTIME_OUTBOUND => $authOutboundBytes,
], $project->getId());
}
}
break;
@@ -997,19 +1047,9 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) {
$projectId = $realtime->connections[$connection]['projectId'];
try {
$consoleDB = getConsoleDB();
$project = $consoleDB->getAuthorization()->skip(
fn () => $consoleDB->getDocument('projects', $projectId)
);
if (!$project->isEmpty()) {
$queue = getQueueForStatsUsageForProject($project);
$queue->addMetric(METRIC_REALTIME_CONNECTIONS, -1)->trigger();
}
} catch (Throwable $th) {
logError($th, 'realtimeUsageConnectionClose', tags: ['projectId' => $projectId]);
}
triggerStats([
METRIC_REALTIME_CONNECTIONS => -1,
], $projectId);
}
$realtime->unsubscribe($connection);