Refactor realtime message handling and enhance query validation tests

This commit is contained in:
ArnabChatterjee20k
2026-04-02 18:32:27 +05:30
parent d8a3b53641
commit bfbf180aee
4 changed files with 135 additions and 102 deletions
+15 -7
View File
@@ -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'];
@@ -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')
);
}
}
@@ -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();
@@ -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();
}