diff --git a/app/init/registers.php b/app/init/registers.php index c07bc9da8b..ceb46c5f56 100644 --- a/app/init/registers.php +++ b/app/init/registers.php @@ -333,7 +333,13 @@ $register->set('pools', function () { $poolAdapter = System::getEnv('_APP_POOL_ADAPTER', default: 'stack') === 'swoole' ? new SwoolePool() : new StackPool(); - $pool = new Pool($poolAdapter, $name, $poolSize, function () use ($type, $resource, $dsn) { + // PubSub workers hold one long-lived subscribed connection and also need + // spare capacity for publishes from the same process. + $connectionPoolSize = $type === 'pubsub' + ? max(2, $poolSize) + : $poolSize; + + $pool = new Pool($poolAdapter, $name, $connectionPoolSize, function () use ($type, $resource, $dsn) { // Get Adapter switch ($type) { case 'database': diff --git a/app/realtime.php b/app/realtime.php index acd15e7540..aae66036d7 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -1,9 +1,9 @@ onMessage(function (int $connection, string $message) use ($server, $re $authorization = new Authorization(); } + // Ensure `$authorization` contains the same roles as the realtime connection. + // `setPermission()` validates against `$authorization->getRoles()`, but roles are + // computed/stored separately in the realtime adapter connection tree. + $connectionRoles = $realtime->connections[$connection]['roles'] ?? []; + foreach ($connectionRoles as $role) { + if ($authorization->hasRole($role)) { + continue; + } + $authorization->addRole($role); + } + $database = getConsoleDB(); $database->setAuthorization($authorization); @@ -1178,6 +1189,10 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re /** @var User $user */ $user = $database->getDocument('users', $userId); + if ($user->isEmpty()) { + throw new Exception(Exception::USER_NOT_FOUND, params:[$userId]); + } + if (!is_array($message['data'])) { throw new Exception(Exception::REALTIME_MESSAGE_FORMAT_INVALID, 'Payload is not valid.'); } diff --git a/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php b/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php index 9889e84105..7ff586b491 100644 --- a/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php +++ b/tests/e2e/Services/Presence/PresenceRealtimeClientTest.php @@ -12,6 +12,7 @@ use Utopia\Database\Helpers\Permission; use Utopia\Database\Helpers\Role; use Utopia\Database\Query; use WebSocket\Client as WebSocketClient; +use WebSocket\TimeoutException; class PresenceRealtimeClientTest extends Scope { @@ -19,6 +20,130 @@ class PresenceRealtimeClientTest extends Scope use RealtimeBase; use SideClient; + private function assertPresenceRealtimeEvent( + array $event, + string $presenceId, + string $action, + string $status, + array $metadata = [], + ?string $expectedUserId = null + ): void { + $expectedUserId ??= $this->getUser()['$id']; + $this->assertSame('event', $event['type'] ?? null); + $this->assertContains('presences', $event['data']['channels'] ?? []); + $this->assertContains('presences.' . $presenceId, $event['data']['channels'] ?? []); + $this->assertNotEmpty($event['data']['events'] ?? []); + $this->assertContains('presences.' . $presenceId . '.' . $action, $event['data']['events'] ?? []); + $this->assertNotEmpty($event['data']['timestamp'] ?? null); + $this->assertArrayHasKey('subscriptions', $event['data'] ?? []); + $this->assertNotEmpty($event['data']['subscriptions'] ?? []); + $this->assertSame($presenceId, $event['data']['payload']['$id'] ?? null); + $this->assertSame($status, $event['data']['payload']['status'] ?? null); + $this->assertSame($metadata, $event['data']['payload']['metadata'] ?? []); + $this->assertSame($expectedUserId, $event['data']['payload']['userId'] ?? null); + } + + private function assertNoRealtimeEvent(WebSocketClient $client): void + { + try { + $client->receive(); + $this->fail('Expected TimeoutException - event should not be received'); + } catch (TimeoutException $e) { + $this->assertTrue(true); + } + } + + /** + * Presence websocket contract: after sending a `type: presence` message, + * the sender socket should receive: + * 1) `type: response` (for the persistence write) + * 2) `type: event` (for the realtime upsert) + * + * This keeps the tests strict about ordering and avoids leaving unread + * realtime events in the socket buffer for later assertions. + */ + private function assertPresenceResponseThenUpsertEvent( + WebSocketClient $client, + string $expectedStatus, + array $expectedMetadata, + ?string $expectedUserId = null + ): string { + $expectedUserId ??= $this->getUser()['$id']; + + $presenceId = null; + $response = null; + $event = null; + + // Ordering is not guaranteed because: + // - response is sent directly via $server->send(...) + // - event is emitted via pub/sub and may arrive earlier + for ($attempts = 0; $attempts < 5; $attempts++) { + $message = \json_decode($client->receive(), true); + $type = $message['type'] ?? null; + + if ($type === 'response') { + $response ??= $message; + + $this->assertSame('presence', $response['data']['to'] ?? null); + $presenceId ??= $response['data']['presence']['$id'] ?? null; + $this->assertNotEmpty($presenceId); + $this->assertSame($expectedStatus, $response['data']['presence']['status'] ?? null); + $this->assertSame($expectedMetadata, $response['data']['presence']['metadata'] ?? null); + } elseif ($type === 'event') { + $event ??= $message; + $presenceId ??= $event['data']['payload']['$id'] ?? null; + $this->assertNotEmpty($presenceId); + + $this->assertPresenceRealtimeEvent( + $event, + $presenceId, + 'upsert', + $expectedStatus, + $expectedMetadata, + $expectedUserId + ); + } + + if ($response !== null && $event !== null) { + return $presenceId; + } + } + + $this->fail('Expected both realtime presence `response` and `event` messages'); + return ''; + } + + /** + * After getting a `response` for a presence message, the next interesting + * realtime message should be the corresponding `event` for the same presence id. + */ + private function receivePresenceEvent( + WebSocketClient $client, + string $presenceId, + string $action, + string $status, + array $expectedMetadata, + ?string $expectedUserId = null + ): array { + do { + $message = \json_decode($client->receive(), true); + } while ( + ($message['type'] ?? null) !== 'event' + || ($message['data']['payload']['$id'] ?? null) !== $presenceId + ); + + $this->assertPresenceRealtimeEvent( + $message, + $presenceId, + $action, + $status, + $expectedMetadata, + $expectedUserId + ); + + return $message; + } + private function connectPresenceSocket(bool $authenticated = true, int $timeout = 2): WebSocketClient { $headers = [ @@ -61,7 +186,7 @@ class PresenceRealtimeClientTest extends Scope { $presenceId = ID::unique(); $userId = $this->getUser()['$id']; - $client = $this->connectPresenceSocket(); + $client = $this->connectPresenceSocket(true, 5); $client->send(\json_encode([ 'type' => 'presence', @@ -75,13 +200,12 @@ class PresenceRealtimeClientTest extends Scope ], ])); - $response = \json_decode($client->receive(), true); - $this->assertSame('response', $response['type'] ?? null); - $this->assertSame('presence', $response['data']['to'] ?? null); - $this->assertSame($presenceId, $response['data']['presence']['$id'] ?? null); - $this->assertSame($userId, $response['data']['presence']['userId'] ?? null); - $this->assertSame('online', $response['data']['presence']['status'] ?? null); - $this->assertSame(['device' => 'web'], $response['data']['presence']['metadata'] ?? null); + $this->assertPresenceResponseThenUpsertEvent( + $client, + 'online', + ['device' => 'web'], + $userId + ); $read = $this->client->call( Client::METHOD_GET, @@ -102,7 +226,7 @@ class PresenceRealtimeClientTest extends Scope { $presenceId = ID::unique(); $userId = $this->getUser()['$id']; - $client = $this->connectPresenceSocket(); + $client = $this->connectPresenceSocket(true, 5); $client->send(\json_encode([ 'type' => 'presence', @@ -115,10 +239,12 @@ class PresenceRealtimeClientTest extends Scope 'permissions' => $this->getPresencePermissions($userId), ], ])); - $first = \json_decode($client->receive(), true); - $this->assertSame('response', $first['type'] ?? null); - $this->assertSame($presenceId, $first['data']['presence']['$id'] ?? null); - $this->assertSame('away', $first['data']['presence']['status'] ?? null); + $this->assertPresenceResponseThenUpsertEvent( + $client, + 'away', + ['source' => 'first'], + $userId + ); $client->send(\json_encode([ 'type' => 'presence', @@ -131,11 +257,12 @@ class PresenceRealtimeClientTest extends Scope 'permissions' => $this->getPresencePermissions($userId), ], ])); - $second = \json_decode($client->receive(), true); - $this->assertSame('response', $second['type'] ?? null); - $this->assertSame($presenceId, $second['data']['presence']['$id'] ?? null); - $this->assertSame('busy', $second['data']['presence']['status'] ?? null); - $this->assertSame(['source' => 'second'], $second['data']['presence']['metadata'] ?? null); + $this->assertPresenceResponseThenUpsertEvent( + $client, + 'busy', + ['source' => 'second'], + $userId + ); $list = $this->client->call( Client::METHOD_GET, @@ -160,6 +287,70 @@ class PresenceRealtimeClientTest extends Scope $client->close(); } + public function testPresenceMessageUpsertWithSameUserPersistsSingleRecord(): void + { + $firstPresenceId = ID::unique(); + $secondPresenceId = ID::unique(); + $userId = $this->getUser()['$id']; + $client = $this->connectPresenceSocket(true, 5); + + $client->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $firstPresenceId, + 'status' => 'away', + 'metadata' => [ + 'source' => 'first-user-upsert', + ], + 'permissions' => $this->getPresencePermissions($userId), + ], + ])); + $this->assertPresenceResponseThenUpsertEvent( + $client, + 'away', + ['source' => 'first-user-upsert'], + $userId + ); + + $client->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $secondPresenceId, + 'status' => 'busy', + 'metadata' => [ + 'source' => 'second-user-upsert', + ], + 'permissions' => $this->getPresencePermissions($userId), + ], + ])); + $this->assertPresenceResponseThenUpsertEvent( + $client, + 'busy', + ['source' => 'second-user-upsert'], + $userId + ); + + $list = $this->client->call( + Client::METHOD_GET, + '/presences', + $this->getServerHeaders(), + [ + 'queries' => [ + Query::equal('userId', [$userId])->toString(), + ], + ] + ); + + $this->assertSame(200, $list['headers']['status-code']); + $this->assertSame(1, $list['body']['total']); + $this->assertCount(1, $list['body']['presences']); + $this->assertSame($userId, $list['body']['presences'][0]['userId']); + $this->assertSame('busy', $list['body']['presences'][0]['status']); + $this->assertSame(['source' => 'second-user-upsert'], $list['body']['presences'][0]['metadata']); + + $client->close(); + } + public function testPresenceMessageValidationErrors(): void { $client = $this->connectPresenceSocket(); @@ -240,13 +431,13 @@ class PresenceRealtimeClientTest extends Scope $this->assertSame(200, $create['headers']['status-code']); $createEvent = \json_decode($client->receive(), true); - $this->assertSame('event', $createEvent['type'] ?? null); - $this->assertContains('presences', $createEvent['data']['channels'] ?? []); - $this->assertContains('presences.' . $presenceId, $createEvent['data']['channels'] ?? []); - $this->assertNotEmpty($createEvent['data']['events'] ?? []); - $this->assertContains('presences.' . $presenceId . '.upsert', $createEvent['data']['events'] ?? []); - $this->assertSame($presenceId, $createEvent['data']['payload']['$id'] ?? null); - $this->assertSame('online', $createEvent['data']['payload']['status'] ?? null); + $this->assertPresenceRealtimeEvent( + $createEvent, + $presenceId, + 'upsert', + 'online', + ['source' => 'channel-parsing-create'] + ); $update = $this->client->call( Client::METHOD_PATCH, @@ -260,13 +451,13 @@ class PresenceRealtimeClientTest extends Scope $this->assertSame(200, $update['headers']['status-code']); $updateEvent = \json_decode($client->receive(), true); - $this->assertSame('event', $updateEvent['type'] ?? null); - $this->assertContains('presences', $updateEvent['data']['channels'] ?? []); - $this->assertContains('presences.' . $presenceId, $updateEvent['data']['channels'] ?? []); - $this->assertNotEmpty($updateEvent['data']['events'] ?? []); - $this->assertContains('presences.' . $presenceId . '.update', $updateEvent['data']['events'] ?? []); - $this->assertSame($presenceId, $updateEvent['data']['payload']['$id'] ?? null); - $this->assertSame('away', $updateEvent['data']['payload']['status'] ?? null); + $this->assertPresenceRealtimeEvent( + $updateEvent, + $presenceId, + 'update', + 'away', + ['source' => 'channel-parsing-update'] + ); $delete = $this->client->call( Client::METHOD_DELETE, @@ -276,12 +467,13 @@ class PresenceRealtimeClientTest extends Scope $this->assertSame(204, $delete['headers']['status-code']); $deleteEvent = \json_decode($client->receive(), true); - $this->assertSame('event', $deleteEvent['type'] ?? null); - $this->assertContains('presences', $deleteEvent['data']['channels'] ?? []); - $this->assertContains('presences.' . $presenceId, $deleteEvent['data']['channels'] ?? []); - $this->assertNotEmpty($deleteEvent['data']['events'] ?? []); - $this->assertContains('presences.' . $presenceId . '.delete', $deleteEvent['data']['events'] ?? []); - $this->assertSame($presenceId, $deleteEvent['data']['payload']['$id'] ?? null); + $this->assertPresenceRealtimeEvent( + $deleteEvent, + $presenceId, + 'delete', + 'away', + ['source' => 'channel-parsing-update'] + ); $client->close(); } @@ -298,6 +490,11 @@ class PresenceRealtimeClientTest extends Scope $listener = $this->getWebsocket(['presences', 'presences.' . $presenceId], $headers, timeout: 8); $connected = \json_decode($listener->receive(), true); $this->assertSame('connected', $connected['type'] ?? null); + $this->assertCount(2, $connected['data']['channels'] ?? []); + $this->assertContains('presences', $connected['data']['channels'] ?? []); + $this->assertContains('presences.' . $presenceId, $connected['data']['channels'] ?? []); + $this->assertCount(1, $connected['data']['subscriptions'] ?? []); + $this->assertNotEmpty(array_values($connected['data']['subscriptions'] ?? [])); $publisher = $this->connectPresenceSocket(true, timeout: 8); @@ -313,24 +510,211 @@ class PresenceRealtimeClientTest extends Scope ], ])); - $createResponse = \json_decode($publisher->receive(), true); - $this->assertSame('response', $createResponse['type'] ?? null); - $this->assertSame('presence', $createResponse['data']['to'] ?? null); - $this->assertSame($presenceId, $createResponse['data']['presence']['$id'] ?? null); + $receivedPresenceId = $this->assertPresenceResponseThenUpsertEvent( + $publisher, + 'online', + ['source' => 'realtime-create-delete-events'], + $userId + ); + $this->assertSame($presenceId, $receivedPresenceId); $createEvent = \json_decode($listener->receive(), true); - $this->assertSame('event', $createEvent['type'] ?? null); - $this->assertContains('presences.' . $presenceId . '.upsert', $createEvent['data']['events'] ?? []); - $this->assertSame($presenceId, $createEvent['data']['payload']['$id'] ?? null); - $this->assertSame('online', $createEvent['data']['payload']['status'] ?? null); + $this->assertPresenceRealtimeEvent( + $createEvent, + $presenceId, + 'upsert', + 'online', + ['source' => 'realtime-create-delete-events'] + ); $publisher->close(); $deleteEvent = \json_decode($listener->receive(), true); - $this->assertSame('event', $deleteEvent['type'] ?? null); - $this->assertContains('presences.' . $presenceId . '.delete', $deleteEvent['data']['events'] ?? []); - $this->assertSame($presenceId, $deleteEvent['data']['payload']['$id'] ?? null); + $this->assertPresenceRealtimeEvent( + $deleteEvent, + $presenceId, + 'delete', + 'online', + ['source' => 'realtime-create-delete-events'] + ); $listener->close(); } + + public function testPresencePermission(): void + { + $presenceIdAny = ID::unique(); + $presenceIdUsers = ID::unique(); + $presenceIdOwner = ID::unique(); + + $user1 = $this->getUser(); + $user1Id = $user1['$id']; + $user2 = $this->getUser(true); + $user3 = $this->getUser(true); + $projectId = $this->getProject()['$id']; + + $user1Headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $user1['session'], + ]; + + $user2Headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $user2['session'], + ]; + + $user3Headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $user3['session'], + ]; + + $user1Listener = $this->getWebsocket(['presences', 'presences.' . $presenceIdAny, 'presences.' . $presenceIdUsers, 'presences.' . $presenceIdOwner], $user1Headers, timeout: 3); + $user2Listener = $this->getWebsocket(['presences', 'presences.' . $presenceIdAny, 'presences.' . $presenceIdUsers, 'presences.' . $presenceIdOwner], $user2Headers, timeout: 3); + $user3Listener = $this->getWebsocket(['presences', 'presences.' . $presenceIdAny, 'presences.' . $presenceIdUsers, 'presences.' . $presenceIdOwner], $user3Headers, timeout: 3); + + $this->assertSame('connected', (\json_decode($user1Listener->receive(), true))['type'] ?? null); + $this->assertSame('connected', (\json_decode($user2Listener->receive(), true))['type'] ?? null); + $this->assertSame('connected', (\json_decode($user3Listener->receive(), true))['type'] ?? null); + + $publisher = $this->getWebsocket(['presences'], $user1Headers, timeout: 5); + $this->assertSame('connected', (\json_decode($publisher->receive(), true))['type'] ?? null); + + $publisher->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $presenceIdAny, + 'status' => 'online', + 'metadata' => [ + 'visibility' => 'any', + ], + 'permissions' => [ + Permission::read(Role::any()), + Permission::update(Role::any()), + Permission::delete(Role::any()), + ], + ], + ])); + + $receivedPresenceId = $this->assertPresenceResponseThenUpsertEvent( + $publisher, + 'online', + ['visibility' => 'any'], + $user1Id + ); + $this->assertSame($presenceIdAny, $receivedPresenceId); + + $this->assertPresenceRealtimeEvent( + \json_decode($user1Listener->receive(), true), + $presenceIdAny, + 'upsert', + 'online', + ['visibility' => 'any'], + $user1Id + ); + $this->assertPresenceRealtimeEvent( + \json_decode($user2Listener->receive(), true), + $presenceIdAny, + 'upsert', + 'online', + ['visibility' => 'any'], + $user1Id + ); + $this->assertPresenceRealtimeEvent( + \json_decode($user3Listener->receive(), true), + $presenceIdAny, + 'upsert', + 'online', + ['visibility' => 'any'], + $user1Id + ); + + $publisher->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $presenceIdUsers, + 'status' => 'away', + 'metadata' => [ + 'visibility' => 'users', + ], + 'permissions' => [ + Permission::read(Role::users()), + Permission::update(Role::users()), + Permission::delete(Role::users()), + ], + ], + ])); + + $receivedPresenceId = $this->assertPresenceResponseThenUpsertEvent( + $publisher, + 'away', + ['visibility' => 'users'], + $user1Id + ); + $this->assertSame($presenceIdUsers, $receivedPresenceId); + + $this->assertPresenceRealtimeEvent( + \json_decode($user1Listener->receive(), true), + $presenceIdUsers, + 'upsert', + 'away', + ['visibility' => 'users'], + $user1Id + ); + $this->assertPresenceRealtimeEvent( + \json_decode($user3Listener->receive(), true), + $presenceIdUsers, + 'upsert', + 'away', + ['visibility' => 'users'], + $user1Id + ); + $this->assertPresenceRealtimeEvent( + \json_decode($user2Listener->receive(), true), + $presenceIdUsers, + 'upsert', + 'away', + ['visibility' => 'users'], + $user1Id + ); + + $publisher->send(\json_encode([ + 'type' => 'presence', + 'data' => [ + 'presenceId' => $presenceIdOwner, + 'status' => 'busy', + 'metadata' => [ + 'visibility' => 'owner', + ], + 'permissions' => [ + Permission::read(Role::user($user1Id)), + Permission::update(Role::user($user1Id)), + Permission::delete(Role::user($user1Id)), + ], + ], + ])); + + $receivedPresenceId = $this->assertPresenceResponseThenUpsertEvent( + $publisher, + 'busy', + ['visibility' => 'owner'], + $user1Id + ); + $this->assertSame($presenceIdOwner, $receivedPresenceId); + + $this->assertPresenceRealtimeEvent( + \json_decode($user1Listener->receive(), true), + $presenceIdOwner, + 'upsert', + 'busy', + ['visibility' => 'owner'], + $user1Id + ); + $this->assertNoRealtimeEvent($user2Listener); + $this->assertNoRealtimeEvent($user3Listener); + + $publisher->close(); + $user1Listener->close(); + $user2Listener->close(); + $user3Listener->close(); + } }