diff --git a/app/realtime.php b/app/realtime.php index 3461ca83e5..71aa251069 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -38,6 +38,7 @@ use Utopia\DSN\DSN; use Utopia\Logger\Log; use Utopia\Pools\Group; use Utopia\Registry\Registry; +use Utopia\Span\Span; use Utopia\System\System; use Utopia\Telemetry\Adapter\None as NoTelemetry; use Utopia\WebSocket\Adapter; @@ -701,6 +702,20 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $project = null; $logUser = null; $authorization = null; + $rawSize = $request->getSize(); + $channelCount = 0; + $subscriptionCount = 0; + $outboundBytes = 0; + $responseCode = 200; + $subscriptionMode = 'message'; + $success = false; + + Span::init('realtime.open'); + Span::add('realtime.connectionId', $connection); + Span::add('realtime.inboundBytes', $rawSize); + if (!empty($request->getOrigin())) { + Span::add('realtime.origin', $request->getOrigin()); + } try { /** @var Document $project */ @@ -750,8 +765,6 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, throw new Exception(Exception::REALTIME_TOO_MANY_MESSAGES, 'Too many requests'); } - $rawSize = $request->getSize(); - triggerStats([ METRIC_REALTIME_INBOUND => $rawSize, ], $project->getId()); @@ -771,6 +784,7 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $roles = $user->getRoles($authorization); $channels = Realtime::convertChannels($request->getQuery('channels', []), $user->getId()); + $channelCount = \count($channels); $updateStats = static function (string $projectId, ?string $teamId, string $payloadJson) use ($register, $stats): void { $register->get('telemetry.connectionCounter')->add(1); @@ -808,11 +822,15 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $realtime->subscribe($project->getId(), $connection, '', $roles, [], [], $user->getId()); $realtime->connections[$connection]['authorization'] = $authorization; $server->send([$connection], $connectedPayloadJson); + $outboundBytes += \strlen($connectedPayloadJson); $updateStats($project->getId(), $project->getAttribute('teamId'), $connectedPayloadJson); + $subscriptionMode = 'message'; + $success = true; return; } $names = array_keys($channels); + $subscriptionMode = 'url'; try { $subscriptions = Realtime::constructSubscriptions( @@ -839,6 +857,7 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $mapping[$index] = $subscriptionId; } + $subscriptionCount = \count($subscriptions); if (!empty($subscriptions)) { $register->get('telemetry.workerSubscriptionCounter')->add(\count($subscriptions), $register->get('telemetry.workerAttributes')); } @@ -857,8 +876,9 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, ]); $server->send([$connection], $connectedPayloadJson); + $outboundBytes += \strlen($connectedPayloadJson); $updateStats($project->getId(), $project->getAttribute('teamId'), $connectedPayloadJson); - + $success = true; } catch (Throwable $th) { logError($th, 'realtime', project: $project, user: $logUser, authorization: $authorization); @@ -868,6 +888,7 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, if (!\is_int($code)) { $code = 500; } + $responseCode = $code; $message = $th->getMessage(); @@ -885,7 +906,9 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, ] ]; - $server->send([$connection], json_encode($response)); + $responsePayloadJson = json_encode($response); + $server->send([$connection], $responsePayloadJson); + $outboundBytes += \strlen($responsePayloadJson); $server->close($connection, $code); if (System::getEnv('_APP_ENV', 'production') === 'development') { @@ -893,16 +916,44 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, Console::error('[Error] Code: ' . $response['data']['code']); Console::error('[Error] Message: ' . $response['data']['message']); } + Span::error($th); + } finally { + Span::add('realtime.success', $success); + Span::add('realtime.responseCode', $responseCode); + Span::add('realtime.subscriptionMode', $subscriptionMode); + Span::add('realtime.channelCount', $channelCount); + Span::add('realtime.subscriptionCount', $subscriptionCount); + Span::add('realtime.outboundBytes', $outboundBytes); + if (!empty($project?->getId())) { + Span::add('realtime.projectId', $project->getId()); + } + if (!empty($logUser?->getId())) { + Span::add('realtime.userId', $logUser->getId()); + } + Span::current()?->finish(); } }); $server->onMessage(function (int $connection, string $message) use ($server, $realtime, $containerId, $register) { $project = null; $authorization = null; + $projectId = $realtime->connections[$connection]['projectId'] ?? null; + $rawSize = \strlen($message); + $messageType = 'invalid'; + $subscriptionDelta = 0; + $subscriptionsRequested = 0; + $subscriptionsRemoved = 0; + $outboundBytes = 0; + $responseCode = 200; + $success = false; + + Span::init('realtime.message'); + Span::add('realtime.connectionId', $connection); + Span::add('realtime.inboundBytes', $rawSize); + Span::add('realtime.containerId', $containerId); + try { - $rawSize = \strlen($message); $response = new Response(new SwooleResponse()); - $projectId = $realtime->connections[$connection]['projectId'] ?? null; // Get authorization from connection (stored during onOpen) $authorization = $realtime->connections[$connection]['authorization'] ?? null; @@ -952,6 +1003,12 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Message format is not valid.'); } + $messageType = $message['type'] ?? 'invalid'; + + if (!\is_scalar($messageType)) { + throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Message type is not valid.'); + } + // Ping does not require project context; other messages do (e.g. after unsubscribe during auth) if (empty($projectId) && ($message['type'] ?? '') !== 'ping') { throw new Exception(Exception::REALTIME_POLICY_VIOLATION, 'Missing project context. Reconnect to the project first.'); @@ -964,6 +1021,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re ]); $server->send([$connection], $pongPayloadJson); + $outboundBytes += \strlen($pongPayloadJson); if ($project !== null && !$project->isEmpty()) { $pongOutboundBytes = \strlen($pongPayloadJson); @@ -1053,6 +1111,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re ]); $server->send([$connection], $authResponsePayloadJson); + $outboundBytes += \strlen($authResponsePayloadJson); if ($project !== null && !$project->isEmpty()) { $authOutboundBytes = \strlen($authResponsePayloadJson); @@ -1114,21 +1173,25 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Invalid query: ' . $e->getMessage()); } + $convertedChannels = \array_keys(Realtime::convertChannels($payload['channels'], $userId)); + $parsedPayloads[] = [ 'subscriptionId' => $subscriptionId, 'channels' => $payload['channels'], + 'convertedChannels' => $convertedChannels, 'queries' => $convertedQueries, ]; } foreach ($parsedPayloads as $parsedPayload) { $subscriptionId = $parsedPayload['subscriptionId']; - $channels = \array_keys(Realtime::convertChannels($parsedPayload['channels'], $userId)); + $channels = $parsedPayload['convertedChannels']; $queries = $parsedPayload['queries']; $realtime->subscribe($projectId, $connection, $subscriptionId, $roles, $channels, $queries); } $subscriptionsAfter = \count($realtime->getSubscriptionMetadata($connection)); $subscriptionDelta = $subscriptionsAfter - $subscriptionsBefore; + $subscriptionsRequested = \count($parsedPayloads); if ($subscriptionDelta !== 0) { $register->get('telemetry.workerSubscriptionCounter')->add($subscriptionDelta, $register->get('telemetry.workerAttributes')); } @@ -1141,7 +1204,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re 'subscriptions' => \array_map(function (array $parsedPayload) { return [ 'subscriptionId' => $parsedPayload['subscriptionId'], - 'channels' => $parsedPayload['channels'], + 'channels' => $parsedPayload['convertedChannels'], 'queries' => \array_map(fn ($q) => $q->toString(), $parsedPayload['queries']), ]; }, $parsedPayloads), @@ -1149,6 +1212,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re ]); $server->send([$connection], $responsePayload); + $outboundBytes += \strlen($responsePayload); if ($project !== null && !$project->isEmpty()) { $subscribeOutboundBytes = \strlen($responsePayload); @@ -1194,6 +1258,8 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re } $subscriptionsAfter = \count($realtime->getSubscriptionMetadata($connection)); $subscriptionDelta = $subscriptionsAfter - $subscriptionsBefore; + $subscriptionsRequested = \count($validatedIds); + $subscriptionsRemoved = \count(\array_filter($unsubscribeResults, fn (array $item) => $item['removed'])); if ($subscriptionDelta !== 0) { $register->get('telemetry.workerSubscriptionCounter')->add($subscriptionDelta, $register->get('telemetry.workerAttributes')); } @@ -1208,6 +1274,7 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re ]); $server->send([$connection], $unsubscribeResponsePayload); + $outboundBytes += \strlen($unsubscribeResponsePayload); if ($project !== null && !$project->isEmpty()) { $unsubscribeOutboundBytes = \strlen($unsubscribeResponsePayload); @@ -1224,12 +1291,14 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re default: throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Message type is not valid.'); } + $success = true; } catch (Throwable $th) { logError($th, 'realtimeMessage', project: $project, authorization: $authorization); $code = $th->getCode(); if (!is_int($code)) { $code = 500; } + $responseCode = $code; $message = $th->getMessage(); @@ -1246,15 +1315,43 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re ] ]; - $server->send([$connection], json_encode($response)); + $responsePayloadJson = json_encode($response); + $server->send([$connection], $responsePayloadJson); + $outboundBytes += \strlen($responsePayloadJson); if ($th->getCode() === 1008) { $server->close($connection, $th->getCode()); } + Span::error($th); + } finally { + Span::add('realtime.success', $success); + Span::add('realtime.responseCode', $responseCode); + Span::add('realtime.subscriptionDelta', $subscriptionDelta); + Span::add('realtime.subscriptionsRequested', $subscriptionsRequested); + Span::add('realtime.subscriptionsRemoved', $subscriptionsRemoved); + Span::add('realtime.subscribe.subscriptionsCount', $subscriptionsRequested); + Span::add('realtime.outboundBytes', $outboundBytes); + Span::add('realtime.projectId', $project?->getId() ?? $projectId); + Span::add('realtime.userId', $realtime->connections[$connection]['userId'] ?? null); + Span::add('realtime.messageType', $messageType); + Span::current()?->finish(); } }); $server->onClose(function (int $connection) use ($realtime, $stats, $register) { + $projectId = null; + $userId = null; + $subscriptionsBeforeClose = 0; + $success = false; + + Span::init('realtime.close'); + Span::add('realtime.connectionId', $connection); + + if (array_key_exists($connection, $realtime->connections)) { + $projectId = $realtime->connections[$connection]['projectId'] ?? null; + $userId = $realtime->connections[$connection]['userId'] ?? null; + } + try { if (array_key_exists($connection, $realtime->connections)) { $stats->decr($realtime->connections[$connection]['projectId'], 'connectionsTotal'); @@ -1271,12 +1368,30 @@ $server->onClose(function (int $connection) use ($realtime, $stats, $register) { METRIC_REALTIME_CONNECTIONS => -1, ], $projectId); } + $success = true; } catch (\Throwable $th) { // Log only; do not rethrow. If we let this bubble, Swoole dumps full coroutine // backtraces and unsubscribe() below would never run (connection cleanup would fail). Console::error('Realtime onClose error: ' . $th->getMessage()); + Span::error($th); + } finally { + try { + $realtime->unsubscribe($connection); + } catch (\Throwable $th) { + Console::error('Realtime onClose unsubscribe error: ' . $th->getMessage()); + Span::error($th); + } + + Span::add('realtime.success', $success); + if (!empty($projectId)) { + Span::add('realtime.projectId', $projectId); + } + if (!empty($userId)) { + Span::add('realtime.userId', $userId); + } + Span::add('realtime.subscriptionsBeforeClose', $subscriptionsBeforeClose); + Span::current()?->finish(); } - $realtime->unsubscribe($connection); Console::info('Connection close: ' . $connection); });