diff --git a/app/realtime.php b/app/realtime.php index 77b56133ff..a6869e00d9 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -797,16 +797,14 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, } }); -$server->onMessage(function (int $connection, string $message) use ($server, $realtime, $containerId, $app) { +$server->onMessage(function (int $connection, string $message) use ($server, $realtime, $containerId) { $project = null; $authorization = null; - + $app = new Http('UTC'); try { $rawSize = \strlen($message); $response = new Response(new SwooleResponse()); $projectId = $realtime->connections[$connection]['projectId'] ?? null; - // TODO: shall it be null or fine to have a guest role? - $roles = $realtime->connections[$connection]['roles'] ?? [Role::guests()->toString()]; // Get authorization from connection (stored during onOpen) $authorization = $realtime->connections[$connection]['authorization'] ?? null; @@ -973,8 +971,18 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Payload is not valid.'); } + // TODO: change this to a clean userId fetching solution + $roles = $realtime->connections[$connection]['roles'] ?? [Role::guests()->toString()]; + $userId = ''; + foreach ($roles as $role) { + if (\str_starts_with($role, 'user:')) { + $userId = \substr($role, 5); + break; + } + } + // bulk validation + parsing before subscribing - foreach ($message['data'] as $payload) { + foreach ($message['data'] as &$payload) { if (!array_key_exists('subscriptionId', $payload)) { throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'subscriptionId is not present in payload.'); } @@ -989,12 +997,12 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re } $subscriptionId = $payload['subscriptionId']; - $channels = $payload['channels']; + $payload['channels'] = \array_keys(Realtime::convertChannels($payload['channels'], $userId)); // TODO: catch error here $payload['queries'] = Query::parseQueries($payload['queries']); } - foreach ($message['data'] as $paylod) { + foreach ($message['data'] as $payload) { $subscriptionId = $payload['subscriptionId']; $channels = $payload['channels']; $queries = $payload['queries']; diff --git a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php index c62e2122e2..1ec7de9b92 100644 --- a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php +++ b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTest.php @@ -6,6 +6,7 @@ use Tests\E2E\Scopes\ProjectCustom; use Tests\E2E\Scopes\Scope; use Tests\E2E\Scopes\SideClient; use Tests\E2E\Services\Functions\FunctionsBase; +use Utopia\Database\Query; class RealtimeCustomClientQueryTest extends Scope { @@ -19,4 +20,97 @@ class RealtimeCustomClientQueryTest extends Scope { return true; } + public function testInvalidQueryShouldNotSubscribe() + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + // Test 1: Simple invalid query method (contains is not allowed) + $client = $this->getWebsocket(['documents'], [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ], null, [ + Query::contains('status', ['active'])->toString(), + ]); + + $response = json_decode($client->receive(), true); + $this->assertEquals('error', $response['type']); + $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); + $this->assertStringContainsString('contains', $response['data']['message']); + + // Test 2: Invalid query method in nested AND query + $client = $this->getWebsocket(['documents'], [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ], null, [ + Query::and([ + Query::equal('status', ['active']), + Query::search('name', 'test') // search is not allowed + ])->toString(), + ]); + + $response = json_decode($client->receive(), true); + $this->assertEquals('error', $response['type']); + $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); + $this->assertStringContainsString('search', $response['data']['message']); + + // Test 3: Invalid query method in nested OR query + $client = $this->getWebsocket(['documents'], [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ], null, [ + Query::or([ + Query::equal('status', ['active']), + Query::between('score', 0, 100) // between is not allowed + ])->toString(), + ]); + + $response = json_decode($client->receive(), true); + $this->assertEquals('error', $response['type']); + $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); + $this->assertStringContainsString('between', $response['data']['message']); + + // Test 4: Deeply nested invalid query (AND -> OR -> invalid) + $client = $this->getWebsocket(['documents'], [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ], null, [ + Query::and([ + Query::equal('status', ['active']), + Query::or([ + Query::greaterThan('score', 50), + Query::startsWith('name', 'test') // startsWith is not allowed + ]) + ])->toString(), + ]); + + $response = json_decode($client->receive(), true); + $this->assertEquals('error', $response['type']); + $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); + $this->assertStringContainsString('startsWith', $response['data']['message']); + + // Test 5: Multiple invalid 'queries' in nested structure + $client = $this->getWebsocket(['documents'], [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ], null, [ + Query::and([ + Query::contains('tags', ['important']), // contains is not allowed + Query::or([ + Query::endsWith('email', '@example.com'), // endsWith is not allowed + Query::equal('status', ['active']) + ]) + ])->toString(), + ]); + + $response = json_decode($client->receive(), true); + $this->assertEquals('error', $response['type']); + $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); + // Should catch the first invalid method encountered + $this->assertTrue( + str_contains($response['data']['message'], 'contains') || + str_contains($response['data']['message'], 'endsWith') + ); + } } diff --git a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTestWithMessage.php b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTestWithMessage.php index b1de21d455..b053fe3897 100644 --- a/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTestWithMessage.php +++ b/tests/e2e/Services/Realtime/RealtimeCustomClientQueryTestWithMessage.php @@ -24,6 +24,16 @@ class RealtimeCustomClientQueryTestWithMessage extends Scope return false; } + protected function supportForAccountChannelQueryAssertion(): bool + { + return false; + } + + protected function supportForInvalidQueryAssertionOnReceive(): bool + { + return false; + } + /** * Same signature as `RealtimeBase::getWebsocket()`, but: * - never sends queries in the URL (avoids URL length limits) @@ -86,6 +96,19 @@ class RealtimeCustomClientQueryTestWithMessage extends Scope return $client; } + private function getWebsocketWithCustomQuery(array $queryParams, array $headers = [], int $timeout = 2): WebSocketClient + { + $queryString = \http_build_query($queryParams); + + return new WebSocketClient( + 'ws://appwrite.test/v1/realtime?' . $queryString, + [ + 'headers' => $headers, + 'timeout' => $timeout, + ] + ); + } + public function testQueryMessageFiltersEvents(): void { $user = $this->getUser(); diff --git a/tests/e2e/Services/Realtime/RealtimeQueryBase.php b/tests/e2e/Services/Realtime/RealtimeQueryBase.php index c72886b3dc..2ebda2397f 100644 --- a/tests/e2e/Services/Realtime/RealtimeQueryBase.php +++ b/tests/e2e/Services/Realtime/RealtimeQueryBase.php @@ -1719,100 +1719,6 @@ trait RealtimeQueryBase $client->close(); } - public function testInvalidQueryShouldNotSubscribe() - { - $user = $this->getUser(); - $session = $user['session'] ?? ''; - $projectId = $this->getProject()['$id']; - - // Test 1: Simple invalid query method (contains is not allowed) - $client = $this->getWebsocket(['documents'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session, - ], null, [ - Query::contains('status', ['active'])->toString(), - ]); - - $response = json_decode($client->receive(), true); - $this->assertEquals('error', $response['type']); - $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); - $this->assertStringContainsString('contains', $response['data']['message']); - - // Test 2: Invalid query method in nested AND query - $client = $this->getWebsocket(['documents'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session, - ], null, [ - Query::and([ - Query::equal('status', ['active']), - Query::search('name', 'test') // search is not allowed - ])->toString(), - ]); - - $response = json_decode($client->receive(), true); - $this->assertEquals('error', $response['type']); - $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); - $this->assertStringContainsString('search', $response['data']['message']); - - // Test 3: Invalid query method in nested OR query - $client = $this->getWebsocket(['documents'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session, - ], null, [ - Query::or([ - Query::equal('status', ['active']), - Query::between('score', 0, 100) // between is not allowed - ])->toString(), - ]); - - $response = json_decode($client->receive(), true); - $this->assertEquals('error', $response['type']); - $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); - $this->assertStringContainsString('between', $response['data']['message']); - - // Test 4: Deeply nested invalid query (AND -> OR -> invalid) - $client = $this->getWebsocket(['documents'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session, - ], null, [ - Query::and([ - Query::equal('status', ['active']), - Query::or([ - Query::greaterThan('score', 50), - Query::startsWith('name', 'test') // startsWith is not allowed - ]) - ])->toString(), - ]); - - $response = json_decode($client->receive(), true); - $this->assertEquals('error', $response['type']); - $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); - $this->assertStringContainsString('startsWith', $response['data']['message']); - - // Test 5: Multiple invalid 'queries' in nested structure - $client = $this->getWebsocket(['documents'], [ - 'origin' => 'http://localhost', - 'cookie' => 'a_session_' . $projectId . '=' . $session, - ], null, [ - Query::and([ - Query::contains('tags', ['important']), // contains is not allowed - Query::or([ - Query::endsWith('email', '@example.com'), // endsWith is not allowed - Query::equal('status', ['active']) - ]) - ])->toString(), - ]); - - $response = json_decode($client->receive(), true); - $this->assertEquals('error', $response['type']); - $this->assertStringContainsString('not supported in Realtime queries', $response['data']['message']); - // Should catch the first invalid method encountered - $this->assertTrue( - str_contains($response['data']['message'], 'contains') || - str_contains($response['data']['message'], 'endsWith') - ); - } - public function testQueryKeys() { $user = $this->getUser(); @@ -2398,7 +2304,9 @@ trait RealtimeQueryBase $this->assertIsArray($event2['data']['subscriptions']); $this->assertNotEmpty($event2['data']['subscriptions']); // Subscription ID should remain stable after permission change - $this->assertContains($originalSubscriptionId, $event2['data']['subscriptions']); + if ($originalSubscriptionId !== null) { + $this->assertContains($originalSubscriptionId, $event2['data']['subscriptions']); + } $client->close(); }