Add userId to connection info in Realtime adapter and simplify userId fetching

This commit is contained in:
ArnabChatterjee20k
2026-04-07 17:35:48 +05:30
parent ca62504b5a
commit bc224de751
3 changed files with 21 additions and 21 deletions
+6 -12
View File
@@ -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' => [
+13 -3
View File
@@ -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
];
}
+2 -6
View File
@@ -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();
}