refactor: remove realtime metrics from project usage endpoints and related tests

This commit is contained in:
ArnabChatterjee20k
2026-03-10 12:15:25 +05:30
parent 0ea196d21c
commit eccc39a466
4 changed files with 1 additions and 367 deletions
-34
View File
@@ -72,10 +72,6 @@ Http::get('/v1/project/usage')
METRIC_DATABASES_OPERATIONS_READS,
METRIC_DATABASES_OPERATIONS_WRITES,
METRIC_FILES_IMAGES_TRANSFORMED,
METRIC_REALTIME_CONNECTIONS,
METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT,
METRIC_REALTIME_INBOUND,
METRIC_REALTIME_OUTBOUND,
],
'period' => [
METRIC_NETWORK_REQUESTS,
@@ -89,10 +85,6 @@ Http::get('/v1/project/usage')
METRIC_DATABASES_OPERATIONS_READS,
METRIC_DATABASES_OPERATIONS_WRITES,
METRIC_FILES_IMAGES_TRANSFORMED,
METRIC_REALTIME_CONNECTIONS,
METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT,
METRIC_REALTIME_INBOUND,
METRIC_REALTIME_OUTBOUND,
]
];
@@ -355,26 +347,6 @@ Http::get('/v1/project/usage')
];
}
// Realtime bandwidth = realtime.inbound + realtime.outbound (per bucket)
$realtimeProjectBandwidth = [];
foreach ($usage[METRIC_REALTIME_INBOUND] as $item) {
$realtimeProjectBandwidth[$item['date']] ??= 0;
$realtimeProjectBandwidth[$item['date']] += $item['value'];
}
foreach ($usage[METRIC_REALTIME_OUTBOUND] as $item) {
$realtimeProjectBandwidth[$item['date']] ??= 0;
$realtimeProjectBandwidth[$item['date']] += $item['value'];
}
$realtimeBandwidth = [];
foreach ($realtimeProjectBandwidth as $date => $value) {
$realtimeBandwidth[] = [
'date' => $date,
'value' => $value
];
}
$response->dynamic(new Document([
'requests' => ($usage[METRIC_NETWORK_REQUESTS]),
'network' => $network,
@@ -395,16 +367,10 @@ Http::get('/v1/project/usage')
'deploymentsStorageTotal' => $total[METRIC_DEPLOYMENTS_STORAGE],
'databasesReadsTotal' => $total[METRIC_DATABASES_OPERATIONS_READS],
'databasesWritesTotal' => $total[METRIC_DATABASES_OPERATIONS_WRITES],
'realtimeConnectionsTotal' => $total[METRIC_REALTIME_CONNECTIONS],
'realtimeMessagesTotal' => $total[METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT],
'realtimeBandwidthTotal' => ($total[METRIC_REALTIME_INBOUND] ?? 0) + ($total[METRIC_REALTIME_OUTBOUND] ?? 0),
'executionsBreakdown' => $executionsBreakdown,
'bucketsBreakdown' => $bucketsBreakdown,
'databasesReads' => $usage[METRIC_DATABASES_OPERATIONS_READS],
'databasesWrites' => $usage[METRIC_DATABASES_OPERATIONS_WRITES],
'realtimeConnections' => $usage[METRIC_REALTIME_CONNECTIONS],
'realtimeMessages' => $usage[METRIC_REALTIME_CONNECTIONS_MESSAGES_SENT],
'realtimeBandwidth' => $realtimeBandwidth,
'databasesStorageBreakdown' => $databasesStorageBreakdown,
'executionsMbSecondsBreakdown' => $executionsMbSecondsBreakdown,
'buildsMbSecondsBreakdown' => $buildsMbSecondsBreakdown,
+1 -58
View File
@@ -1,6 +1,5 @@
<?php
use Appwrite\Event\StatsUsage;
use Appwrite\Extend\Exception;
use Appwrite\Extend\Exception as AppwriteException;
use Appwrite\Messaging\Adapter\Realtime;
@@ -38,7 +37,6 @@ use Utopia\DSN\DSN;
use Utopia\Http\Http;
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;
@@ -226,65 +224,10 @@ if (!function_exists('getTelemetry')) {
}
}
if (!function_exists('getQueueForStatsUsageForProject')) {
function getQueueForStatsUsageForProject(Document $project): StatsUsage
{
$ctx = Coroutine::getContext();
if (!isset($ctx['queueForStatsUsage'])) {
$ctx['queueForStatsUsage'] = [];
}
if (isset($ctx['queueForStatsUsage'][$project->getSequence()])) {
return $ctx['queueForStatsUsage'][$project->getSequence()];
}
global $register;
/** @var Group $pools */
$pools = $register->get('pools');
$queue = new StatsUsage(new BrokerPool(publisher: $pools->get('publisher')));
$queue->setProject($project);
return $ctx['queueForStatsUsage'][$project->getSequence()] = $queue;
}
}
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]);
}
return;
}
}
@@ -100,24 +100,6 @@ class UsageProject extends Model
'default' => 0,
'example' => 0,
])
->addRule('realtimeConnectionsTotal', [
'type' => self::TYPE_INTEGER,
'description' => 'Current aggregated number of open Realtime connections.',
'default' => 0,
'example' => 0,
])
->addRule('realtimeMessagesTotal', [
'type' => self::TYPE_INTEGER,
'description' => 'Total number of Realtime messages sent to clients.',
'default' => 0,
'example' => 0,
])
->addRule('realtimeBandwidthTotal', [
'type' => self::TYPE_INTEGER,
'description' => 'Total consumed Realtime bandwidth (in bytes).',
'default' => 0,
'example' => 0,
])
->addRule('requests', [
'type' => Response::MODEL_METRIC,
'description' => 'Aggregated number of requests per period.',
@@ -132,27 +114,6 @@ class UsageProject extends Model
'example' => [],
'array' => true
])
->addRule('realtimeConnections', [
'type' => Response::MODEL_METRIC,
'description' => 'Aggregated number of open Realtime connections per period.',
'default' => [],
'example' => [],
'array' => true
])
->addRule('realtimeMessages', [
'type' => Response::MODEL_METRIC,
'description' => 'Aggregated number of Realtime messages sent to clients per period.',
'default' => [],
'example' => [],
'array' => true
])
->addRule('realtimeBandwidth', [
'type' => Response::MODEL_METRIC,
'description' => 'Aggregated consumed Realtime bandwidth (in bytes) per period.',
'default' => [],
'example' => [],
'array' => true
])
->addRule('users', [
'type' => Response::MODEL_METRIC,
'description' => 'Aggregated number of users per period.',
-236
View File
@@ -1410,242 +1410,6 @@ class UsageTest extends Scope
});
}
public function testRealtimeUsageMetrics(): void
{
$user = $this->getUser();
$session = $user['session'] ?? '';
$projectId = $this->getProject()['$id'];
// Baseline realtime usage before opening a new connection
$baseline = $this->client->call(
Client::METHOD_GET,
'/project/usage',
$this->getConsoleHeaders(),
[
'period' => '1h',
'startDate' => self::getToday(),
'endDate' => self::getTomorrow(),
]
);
$connectionsBefore = $baseline['body']['realtimeConnectionsTotal'] ?? 0;
$messagesBefore = $baseline['body']['realtimeMessagesTotal'] ?? 0;
$connectionCount = 3;
$clients = [];
for ($i = 0; $i < $connectionCount; $i++) {
$client = $this->getWebsocket(['documents'], [
'origin' => 'http://localhost',
'cookie' => 'a_session_' . $projectId . '=' . $session,
], null);
$connected = json_decode($client->receive(), true);
$this->assertEquals('connected', $connected['type']);
$clients[] = $client;
}
try {
$database = $this->client->call(Client::METHOD_POST, '/databases', [
'content-type' => 'application/json',
'x-appwrite-project' => $projectId,
'x-appwrite-key' => $this->getProject()['apiKey'],
], [
'databaseId' => ID::unique(),
'name' => 'Realtime Usage DB',
]);
$databaseId = $database['body']['$id'];
$collection = $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections', [
'content-type' => 'application/json',
'x-appwrite-project' => $projectId,
'x-appwrite-key' => $this->getProject()['apiKey'],
], [
'collectionId' => ID::unique(),
'name' => 'Realtime Usage Collection',
'permissions' => [
Permission::create(Role::user($user['$id'])),
],
'documentSecurity' => true,
]);
$collectionId = $collection['body']['$id'];
$attribute = $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections/' . $collectionId . '/attributes/string', [
'content-type' => 'application/json',
'x-appwrite-project' => $projectId,
'x-appwrite-key' => $this->getProject()['apiKey'],
], [
'key' => 'name',
'size' => 256,
'required' => true,
]);
$this->assertEquals(202, $attribute['headers']['status-code']);
$this->assertEventually(function () use ($databaseId, $collectionId, $projectId) {
$response = $this->client->call(Client::METHOD_GET, '/databases/' . $databaseId . '/collections/' . $collectionId . '/attributes/name', [
'content-type' => 'application/json',
'x-appwrite-project' => $projectId,
'x-appwrite-key' => $this->getProject()['apiKey'],
]);
$this->assertEquals('available', $response['body']['status']);
}, 30000, 250);
$document = $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections/' . $collectionId . '/documents', array_merge([
'content-type' => 'application/json',
'x-appwrite-project' => $projectId,
], $this->getHeaders()), [
'documentId' => ID::unique(),
'data' => [
'name' => 'Realtime Usage Doc',
],
'permissions' => [
Permission::read(Role::any()),
Permission::update(Role::any()),
Permission::delete(Role::any()),
],
]);
$this->assertEquals(201, $document['headers']['status-code']);
$event = json_decode($clients[0]->receive(), true);
$this->assertEquals('event', $event['type']);
// After creating a document we expect all connections to receive an event
$this->assertEventually(function () use ($connectionsBefore, $messagesBefore, $connectionCount) {
$response = $this->client->call(
Client::METHOD_GET,
'/project/usage',
$this->getConsoleHeaders(),
[
'period' => '1h',
'startDate' => self::getToday(),
'endDate' => self::getTomorrow(),
]
);
$this->assertEquals(200, $response['headers']['status-code']);
$this->assertArrayHasKey('realtimeConnectionsTotal', $response['body']);
$this->assertArrayHasKey('realtimeMessagesTotal', $response['body']);
$this->assertArrayHasKey('realtimeBandwidthTotal', $response['body']);
$this->assertArrayHasKey('realtimeConnections', $response['body']);
$this->assertArrayHasKey('realtimeMessages', $response['body']);
$this->assertArrayHasKey('realtimeBandwidth', $response['body']);
// We expect exactly $connectionCount additional open connections and $connectionCount additional message deliveries
$this->assertEquals($connectionsBefore + $connectionCount, $response['body']['realtimeConnectionsTotal']);
$this->assertEquals($messagesBefore + $connectionCount, $response['body']['realtimeMessagesTotal']);
$this->assertGreaterThan(0, $response['body']['realtimeBandwidthTotal']);
$this->validateDates($response['body']['realtimeConnections']);
$this->validateDates($response['body']['realtimeMessages']);
$this->validateDates($response['body']['realtimeBandwidth']);
}, 60000, 2000);
// Capture a snapshot of usage after the broadcasted document event
$afterEventUsage = $this->client->call(
Client::METHOD_GET,
'/project/usage',
$this->getConsoleHeaders(),
[
'period' => '1h',
'startDate' => self::getToday(),
'endDate' => self::getTomorrow(),
]
);
$this->assertEquals(200, $afterEventUsage['headers']['status-code']);
$connectionsAfterEvent = $afterEventUsage['body']['realtimeConnectionsTotal'] ?? 0;
$messagesAfterEvent = $afterEventUsage['body']['realtimeMessagesTotal'] ?? 0;
$bandwidthAfterEvent = $afterEventUsage['body']['realtimeBandwidthTotal'] ?? 0;
// Send a ping over an existing connection to exercise the ping/pong
// realtime usage metrics path (inbound + outbound bytes) without
// generating additional "messages sent" usage.
$clients[1]->send(json_encode([
'type' => 'ping',
]));
// A broadcast document "event" can still be queued for this connection.
// Read until we see the pong response (and fail fast on unexpected frames).
$pong = null;
for ($i = 0; $i < 5; $i++) {
$frame = json_decode($clients[1]->receive(), true);
$this->assertIsArray($frame);
$this->assertArrayHasKey('type', $frame);
if ($frame['type'] === 'pong') {
$pong = $frame;
break;
}
$this->assertEquals('event', $frame['type']);
}
$this->assertNotNull($pong, 'Expected to receive a pong frame after ping.');
// We expect:
// - connections count to remain the same
// - messages count to remain the same (no new broadcast events)
// - bandwidth total to increase because of ping/pong traffic
$this->assertEventually(function () use ($connectionsAfterEvent, $messagesAfterEvent, $bandwidthAfterEvent) {
$response = $this->client->call(
Client::METHOD_GET,
'/project/usage',
$this->getConsoleHeaders(),
[
'period' => '1h',
'startDate' => self::getToday(),
'endDate' => self::getTomorrow(),
]
);
$this->assertEquals(200, $response['headers']['status-code']);
$this->assertEquals($connectionsAfterEvent, $response['body']['realtimeConnectionsTotal']);
$this->assertEquals($messagesAfterEvent, $response['body']['realtimeMessagesTotal']);
$this->assertGreaterThan($bandwidthAfterEvent, $response['body']['realtimeBandwidthTotal']);
}, 60000, 2000);
// Now close a single connection and ensure the counters reflect it
$clients[0]->close();
$this->assertEventually(function () use ($connectionsBefore, $messagesBefore, $connectionCount) {
$response = $this->client->call(
Client::METHOD_GET,
'/project/usage',
$this->getConsoleHeaders(),
[
'period' => '1h',
'startDate' => self::getToday(),
'endDate' => self::getTomorrow(),
]
);
$this->assertEquals(200, $response['headers']['status-code']);
$this->assertArrayHasKey('realtimeConnectionsTotal', $response['body']);
$this->assertArrayHasKey('realtimeMessagesTotal', $response['body']);
$this->assertArrayHasKey('realtimeBandwidthTotal', $response['body']);
// One of the connections is closed, so we expect one less open connection.
// Messages and bandwidth are cumulative and should not decrease.
$this->assertEquals($connectionsBefore + $connectionCount - 1, $response['body']['realtimeConnectionsTotal']);
$this->assertEquals($messagesBefore + $connectionCount, $response['body']['realtimeMessagesTotal']);
$this->assertGreaterThan(0, $response['body']['realtimeBandwidthTotal']);
}, 60000, 2000);
} finally {
foreach ($clients as $client) {
$client->close();
}
}
}
public function tearDown(): void
{
$this->projectId = '';