diff --git a/app/realtime.php b/app/realtime.php index 3bceb54f26..52bf62a11f 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -711,7 +711,7 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, ] ]); - $realtime->subscribe($project->getId(), $connection, '', $roles, [], []); + $realtime->subscribe($project->getId(), $connection, '', $roles, [], [], $user->getId()); $server->send([$connection], $connectedPayloadJson); return; } @@ -737,7 +737,8 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $subscriptionId, $roles, $subscription['channels'], - $subscription['queries'] + $subscription['queries'], + $user->getId() ); $mapping[$index] = $subscriptionId; @@ -937,7 +938,8 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re $subscriptionId, $roles, $subscription['channels'] ?? [], - $queries + $queries, + $user->getId() ); } } @@ -985,15 +987,8 @@ $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; - } - } + $userId = $realtime->connections[$connection]['userId'] ?? ''; // bulk validation + parsing before subscribing $parsedPayloads = []; @@ -1039,7 +1034,6 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re $realtime->subscribe($projectId, $connection, $subscriptionId, $roles, $channels, $queries); } - // TODO: find a better way to store the queries and no reconversion $responsePayload = json_encode([ 'type' => 'response', 'data' => [ diff --git a/src/Appwrite/Messaging/Adapter/Realtime.php b/src/Appwrite/Messaging/Adapter/Realtime.php index f9e97f0e45..f1d806bcc5 100644 --- a/src/Appwrite/Messaging/Adapter/Realtime.php +++ b/src/Appwrite/Messaging/Adapter/Realtime.php @@ -20,6 +20,7 @@ class Realtime extends MessagingAdapter * [CONNECTION_ID] -> * 'projectId' -> [PROJECT_ID] * 'roles' -> [ROLE_x, ROLE_Y] + * 'userId' -> [USER_ID] * 'channels' -> [CHANNEL_NAME_X, CHANNEL_NAME_Y, CHANNEL_NAME_Z] */ public array $connections = []; @@ -67,8 +68,15 @@ class Realtime extends MessagingAdapter * @param array $queryGroup Array of Query objects for this subscription (AND logic within subscription) * @return void */ - public function subscribe(string $projectId, mixed $identifier, string $subscriptionId, array $roles, array $channels, array $queryGroup = []): void - { + public function subscribe( + string $projectId, + mixed $identifier, + string $subscriptionId, + array $roles, + array $channels, + array $queryGroup = [], + ?string $userId = null + ): void { if (!isset($this->subscriptions[$projectId])) { // Init Project $this->subscriptions[$projectId] = []; } @@ -106,10 +114,12 @@ class Realtime extends MessagingAdapter } } - // Update connection info + // Keep userId from onOpen/authentication when provided. + // Fallback to existing stored value for subsequent subscribe upserts. $this->connections[$identifier] = [ 'projectId' => $projectId, 'roles' => $roles, + 'userId' => $userId ?? ($this->connections[$identifier]['userId'] ?? ''), 'channels' => $channels ]; } diff --git a/tests/e2e/Services/Realtime/RealtimeBase.php b/tests/e2e/Services/Realtime/RealtimeBase.php index 95f3665e4c..b2d17c1e4a 100644 --- a/tests/e2e/Services/Realtime/RealtimeBase.php +++ b/tests/e2e/Services/Realtime/RealtimeBase.php @@ -101,18 +101,14 @@ trait RealtimeBase $client->close(); } - public function testConnectionFailureMissingChannels(): void + public function testConnectionSuccessMissingChannels(): void { $client = $this->getWebsocket([]); $payload = json_decode($client->receive(), true); $this->assertArrayHasKey("type", $payload); $this->assertArrayHasKey("data", $payload); - $this->assertEquals("error", $payload["type"]); - $this->assertEquals(1008, $payload["data"]["code"]); - $this->assertEquals("Missing channels", $payload["data"]["message"]); - \usleep(250000); // 250ms - $this->expectException(ConnectionException::class); // Check if server disconnected client + $this->assertEquals("connected", $payload["type"]); $client->close(); }