From d2423a5bb51e5fe0d1ec34c9c38458caa77fefbd Mon Sep 17 00:00:00 2001 From: ArnabChatterjee20k Date: Mon, 27 Apr 2026 12:57:35 +0530 Subject: [PATCH] added tests --- .../Services/Realtime/RealtimeQueryBase.php | 420 ++++++++++++++++++ 1 file changed, 420 insertions(+) diff --git a/tests/e2e/Services/Realtime/RealtimeQueryBase.php b/tests/e2e/Services/Realtime/RealtimeQueryBase.php index 04b8400b57..24d2a3511a 100644 --- a/tests/e2e/Services/Realtime/RealtimeQueryBase.php +++ b/tests/e2e/Services/Realtime/RealtimeQueryBase.php @@ -2446,4 +2446,424 @@ trait RealtimeQueryBase $clientWithMatchingQuery->close(); $clientWithNonMatchingQuery->close(); } + + /** + * Sets up a database + collection + 'name' string attribute, returning their IDs. + * Used by action-channel tests to avoid duplicating fixture code. + * + * @return array{databaseId: string, collectionId: string} + */ + private function createActorsCollection(): array + { + $database = $this->client->call(Client::METHOD_POST, '/databases', array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $this->getProject()['$id'], + 'x-appwrite-key' => $this->getProject()['apiKey'], + ]), [ + 'databaseId' => ID::unique(), + 'name' => 'Action Channel DB', + ]); + $databaseId = $database['body']['$id']; + + $collection = $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections', array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $this->getProject()['$id'], + 'x-appwrite-key' => $this->getProject()['apiKey'], + ]), [ + 'collectionId' => ID::unique(), + 'name' => 'Actors', + 'permissions' => [ + Permission::create(Role::user($this->getUser()['$id'])), + ], + 'documentSecurity' => true, + ]); + $collectionId = $collection['body']['$id']; + + $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections/' . $collectionId . '/attributes/string', array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $this->getProject()['$id'], + 'x-appwrite-key' => $this->getProject()['apiKey'], + ]), [ + 'key' => 'name', + 'size' => 256, + 'required' => true, + ]); + + $this->assertEventually(function () use ($databaseId, $collectionId) { + $response = $this->client->call(Client::METHOD_GET, '/databases/' . $databaseId . '/collections/' . $collectionId . '/attributes/name', array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $this->getProject()['$id'], + 'x-appwrite-key' => $this->getProject()['apiKey'], + ])); + $this->assertEquals('available', $response['body']['status']); + }, 30000, 250); + + return ['databaseId' => $databaseId, 'collectionId' => $collectionId]; + } + + /** + * Creates a document with the given ID and name. Returns the parsed body. + * Permissions allow Role::any() for all CRUD so any session can observe the events. + * + * @return array + */ + private function createActor(string $databaseId, string $collectionId, string $documentId, string $name): array + { + $document = $this->client->call(Client::METHOD_POST, '/databases/' . $databaseId . '/collections/' . $collectionId . '/documents', array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $this->getProject()['$id'], + ], $this->getHeaders()), [ + 'documentId' => $documentId, + 'data' => ['name' => $name], + 'permissions' => [ + Permission::read(Role::any()), + Permission::update(Role::any()), + Permission::delete(Role::any()), + ], + ]); + + return $document['body']; + } + + public function testChannelActionFilterReflectedInConnectedResponse(): void + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ]; + + // Subscribing with an action suffix should round-trip the original channel + // name on the connected response. Only meaningful in URL-subscribe mode — + // the message-based path consumes the connected response inside its + // getWebsocket helper before returning, so we can't observe it here. + $client = $this->getWebsocket([ + 'documents.create', + 'documents.update', + 'documents.upsert', + 'documents', + ], $headers); + + $connected = $this->assertConnectionStatusIfSupported($client); + if ($connected === null) { + $client->close(); + $this->markTestSkipped('Connected-response channels are not surfaced through the message-based subscribe path.'); + } + + $this->assertContains('documents.create', $connected['data']['channels']); + $this->assertContains('documents.update', $connected['data']['channels']); + $this->assertContains('documents.upsert', $connected['data']['channels']); + $this->assertContains('documents', $connected['data']['channels']); + + $client->close(); + } + + public function testChannelActionFilterDeliversOnlyMatchingActions(): void + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ]; + + ['databaseId' => $databaseId, 'collectionId' => $collectionId] = $this->createActorsCollection(); + + $createChannel = "databases.{$databaseId}.collections.{$collectionId}.documents.create"; + $updateChannel = "databases.{$databaseId}.collections.{$collectionId}.documents.update"; + $upsertChannel = "databases.{$databaseId}.collections.{$collectionId}.documents.upsert"; + + $clientCreate = $this->getWebsocket([$createChannel], $headers); + $clientUpdate = $this->getWebsocket([$updateChannel], $headers); + $clientUpsert = $this->getWebsocket([$upsertChannel], $headers); + + $this->assertConnectionStatusIfSupported($clientCreate); + $this->assertConnectionStatusIfSupported($clientUpdate); + $this->assertConnectionStatusIfSupported($clientUpsert); + + $documentId = ID::unique(); + $this->createActor($databaseId, $collectionId, $documentId, 'Chris Evans'); + + // Create event delivers only to the .create subscriber. + $createEvent = json_decode($clientCreate->receive(), true); + $this->assertEquals('event', $createEvent['type']); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$documentId}.create", + $createEvent['data']['events'] + ); + $this->assertEquals('Chris Evans', $createEvent['data']['payload']['name']); + + try { + $clientUpdate->receive(); + $this->fail('Update subscriber should not receive a create event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + try { + $clientUpsert->receive(); + $this->fail('Upsert subscriber should not receive a create event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + // Update fires update events; only the .update subscriber should hear them. + $this->client->call(Client::METHOD_PATCH, "/databases/{$databaseId}/collections/{$collectionId}/documents/{$documentId}", array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $projectId, + ], $this->getHeaders()), [ + 'data' => ['name' => 'Chris Evans 2'], + ]); + + $updateEvent = json_decode($clientUpdate->receive(), true); + $this->assertEquals('event', $updateEvent['type']); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$documentId}.update", + $updateEvent['data']['events'] + ); + $this->assertEquals('Chris Evans 2', $updateEvent['data']['payload']['name']); + + try { + $clientCreate->receive(); + $this->fail('Create subscriber should not receive an update event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + try { + $clientUpsert->receive(); + $this->fail('Upsert subscriber should not receive an update event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + // PUT bulk upsert fires upsert events; only the .upsert subscriber should hear them. + $this->client->call(Client::METHOD_PUT, "/databases/{$databaseId}/collections/{$collectionId}/documents", array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $projectId, + 'x-appwrite-key' => $this->getProject()['apiKey'], + ]), [ + 'documents' => [ + [ + '$id' => ID::unique(), + 'name' => 'Robert Downey Jr.', + '$permissions' => [ + Permission::read(Role::any()), + Permission::update(Role::any()), + Permission::delete(Role::any()), + ], + ], + ], + ]); + + $upsertEvent = json_decode($clientUpsert->receive(), true); + $this->assertEquals('event', $upsertEvent['type']); + $this->assertContains( + "databases.{$databaseId}.collections.*.documents.*.upsert", + $upsertEvent['data']['events'] + ); + + try { + $clientCreate->receive(); + $this->fail('Create subscriber should not receive an upsert event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + try { + $clientUpdate->receive(); + $this->fail('Update subscriber should not receive an upsert event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + $clientCreate->close(); + $clientUpdate->close(); + $clientUpsert->close(); + } + + public function testChannelActionFilterByDocumentId(): void + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ]; + + ['databaseId' => $databaseId, 'collectionId' => $collectionId] = $this->createActorsCollection(); + + // Use a known custom ID so the .id.action channel can be subscribed before the + // document exists. Without this the channel name can't be predicted. + $watchedId = 'actor-watched'; + $idCreateChannel = "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}.create"; + + $clientWatched = $this->getWebsocket([$idCreateChannel], $headers); + $connected = $this->assertConnectionStatusIfSupported($clientWatched); + if ($connected !== null) { + $this->assertContains($idCreateChannel, $connected['data']['channels']); + } + + // Creating a *different* document should not trigger the watched-id subscription. + $this->createActor($databaseId, $collectionId, ID::unique(), 'Other Actor'); + + try { + $clientWatched->receive(); + $this->fail('Subscriber to .{id}.create should not receive events for a different document.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + // Creating the watched document delivers exactly one create event. + $this->createActor($databaseId, $collectionId, $watchedId, 'Watched Actor'); + + $event = json_decode($clientWatched->receive(), true); + $this->assertEquals('event', $event['type']); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}.create", + $event['data']['events'] + ); + $this->assertEquals($watchedId, $event['data']['payload']['$id']); + $this->assertEquals('Watched Actor', $event['data']['payload']['name']); + + // Updating the watched document does NOT match — action filter is `create` only. + $this->client->call(Client::METHOD_PATCH, "/databases/{$databaseId}/collections/{$collectionId}/documents/{$watchedId}", array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $projectId, + ], $this->getHeaders()), [ + 'data' => ['name' => 'Watched Actor v2'], + ]); + + try { + $clientWatched->receive(); + $this->fail('Subscriber to .{id}.create should not receive update events on the same document.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + $clientWatched->close(); + } + + public function testChannelActionFilterMultiChannelSubscription(): void + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ]; + + ['databaseId' => $databaseId, 'collectionId' => $collectionId] = $this->createActorsCollection(); + + $watchedId = 'actor-multi'; + $idCreateChannel = "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}.create"; + $rowsChannel = "databases.{$databaseId}.tables.{$collectionId}.rows"; + + // One subscription that listens on both: + // 1. `databases...documents.{watchedId}.create` — narrow, action-filtered + // 2. `databases...tables.{collectionId}.rows` — broad, non-action (tablesdb mirror) + // A create on the watched document must reach this subscriber via *both* channels. + $clientMulti = $this->getWebsocket([$idCreateChannel, $rowsChannel], $headers); + $connected = $this->assertConnectionStatusIfSupported($clientMulti); + if ($connected !== null) { + $this->assertContains($idCreateChannel, $connected['data']['channels']); + $this->assertContains($rowsChannel, $connected['data']['channels']); + } + + $this->createActor($databaseId, $collectionId, $watchedId, 'Multi Actor'); + + $event = json_decode($clientMulti->receive(), true); + $this->assertEquals('event', $event['type']); + // The event payload's channels list reports the underlying base channels that + // the published event carries. Both the broad rows channel and the document + // channel that the action filter is anchored on should be present. + $this->assertContains($rowsChannel, $event['data']['channels']); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}", + $event['data']['channels'] + ); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}.create", + $event['data']['events'] + ); + $this->assertEquals('Multi Actor', $event['data']['payload']['name']); + + // Update on the same doc: the .{id}.create branch is filtered out, but the + // broad rows channel has no action filter — the subscription still receives + // the event via that branch (a single delivery, not two). + $this->client->call(Client::METHOD_PATCH, "/databases/{$databaseId}/collections/{$collectionId}/documents/{$watchedId}", array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $projectId, + ], $this->getHeaders()), [ + 'data' => ['name' => 'Multi Actor v2'], + ]); + + $update = json_decode($clientMulti->receive(), true); + $this->assertEquals('event', $update['type']); + $this->assertContains($rowsChannel, $update['data']['channels']); + $this->assertContains( + "databases.{$databaseId}.collections.{$collectionId}.documents.{$watchedId}.update", + $update['data']['events'] + ); + + // No second copy of the same update should arrive — getSubscribers folds + // multi-channel matches into a single connection delivery. + try { + $clientMulti->receive(); + $this->fail('Multi-channel subscriber should receive a single delivery per event.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + $clientMulti->close(); + } + + public function testChannelActionFilterUnsupportedActionTreatedAsLiteral(): void + { + $user = $this->getUser(); + $session = $user['session'] ?? ''; + $projectId = $this->getProject()['$id']; + + $headers = [ + 'origin' => 'http://localhost', + 'cookie' => 'a_session_' . $projectId . '=' . $session, + ]; + + ['databaseId' => $databaseId, 'collectionId' => $collectionId] = $this->createActorsCollection(); + + // `delete` is intentionally NOT in SUPPORTED_ACTIONS yet, so parseActionChannel + // leaves the channel name intact and treats it as a literal channel that no + // published event ever carries — the subscriber should receive nothing. + $client = $this->getWebsocket(['documents.delete'], $headers); + $connected = $this->assertConnectionStatusIfSupported($client); + if ($connected !== null) { + $this->assertContains('documents.delete', $connected['data']['channels']); + } + + $documentId = ID::unique(); + $this->createActor($databaseId, $collectionId, $documentId, 'No Delete Listener'); + + $this->client->call(Client::METHOD_DELETE, "/databases/{$databaseId}/collections/{$collectionId}/documents/{$documentId}", array_merge([ + 'content-type' => 'application/json', + 'x-appwrite-project' => $projectId, + ], $this->getHeaders())); + + try { + $client->receive(); + $this->fail('`documents.delete` is not (yet) a supported action channel and should not deliver.'); + } catch (TimeoutException $e) { + $this->addToAssertionCount(1); + } + + $client->close(); + } }