diff --git a/src/Appwrite/Messaging/Adapter/Realtime.php b/src/Appwrite/Messaging/Adapter/Realtime.php index 8fe7342ec2..a03ccecfd9 100644 --- a/src/Appwrite/Messaging/Adapter/Realtime.php +++ b/src/Appwrite/Messaging/Adapter/Realtime.php @@ -14,6 +14,18 @@ use Utopia\Database\Query; class Realtime extends MessagingAdapter { + /** + * Action suffix that means "all actions" — i.e. no action filter on this subscription. + */ + public const ACTION_ALL = '*'; + + /** + * Action suffixes recognized in channel names. A channel like `documents.create` + * is split into base channel `documents` plus action `create`. Add new actions + * (e.g. `delete`) here to extend support — no other code change is required. + */ + public const SUPPORTED_ACTIONS = ['create', 'update', 'upsert']; + /** * Connection Tree * @@ -21,7 +33,11 @@ class Realtime extends MessagingAdapter * 'projectId' -> [PROJECT_ID] * 'roles' -> [ROLE_x, ROLE_Y] * 'userId' -> [USER_ID] - * 'channels' -> [CHANNEL_NAME_X, CHANNEL_NAME_Y, CHANNEL_NAME_Z] + * 'channels' -> [BASE_CHANNEL_X, BASE_CHANNEL_Y, BASE_CHANNEL_Z] + * + * Channels here are stored in their *base* form (action suffix stripped) so they + * line up with subscription-tree keys; the original action-prefixed channel is + * reconstructed from per-subscription `actions` metadata when needed. */ public array $connections = []; @@ -32,9 +48,13 @@ class Realtime extends MessagingAdapter * [ROLE_X] -> * [CHANNEL_NAME_X] -> * [CONNECTION_ID] -> - * [SUB_ID] -> ['strings' => [...], 'compiled' => [...]] + * [SUB_ID] -> ['strings' => [...], 'compiled' => [...], 'actions' => [...]] + * + * Each subscription ID maps to query strings (for metadata), pre-compiled query + * filters, and an `actions` metadata list. `actions` is `['*']` by default + * meaning "no action filter"; otherwise a list of action suffixes (e.g. `['create']`) + * that the event must end with for delivery. * - * Each subscription ID maps to query strings (for metadata) and pre-compiled query filters. * Within a subscription: AND logic (all queries must match) * Across subscriptions: OR logic (any subscription matching = send event) */ @@ -81,10 +101,33 @@ class Realtime extends MessagingAdapter $this->subscriptions[$projectId] = []; } + // Split each channel into a base form + an action suffix. Channels without a + // recognised suffix get action '*' (no filter). When the same base channel + // appears with multiple actions in this call (e.g. `documents.create` and + // `documents.update`), the actions are merged onto a single tree entry. + $actionsByBase = []; + $baseChannels = []; + foreach ($channels as $channel) { + [$base, $action] = self::parseActionChannel($channel); + if (!\in_array($base, $baseChannels, true)) { + $baseChannels[] = $base; + } + if (!isset($actionsByBase[$base])) { + $actionsByBase[$base] = [$action]; + continue; + } + // '*' subsumes any specific action — once present, drop the rest. + if (\in_array(self::ACTION_ALL, $actionsByBase[$base], true) || $action === self::ACTION_ALL) { + $actionsByBase[$base] = [self::ACTION_ALL]; + } elseif (!\in_array($action, $actionsByBase[$base], true)) { + $actionsByBase[$base][] = $action; + } + } + $strings = []; $data = []; - if (!empty($channels)) { + if (!empty($baseChannels)) { if (empty($queryGroup)) { $strings[] = Query::select(['*'])->toString(); } else { @@ -103,19 +146,24 @@ class Realtime extends MessagingAdapter $this->subscriptions[$projectId][$role] = []; } - foreach ($channels as $channel) { - if (!isset($this->subscriptions[$projectId][$role][$channel])) { - $this->subscriptions[$projectId][$role][$channel] = []; + foreach ($actionsByBase as $base => $actions) { + if (!isset($this->subscriptions[$projectId][$role][$base])) { + $this->subscriptions[$projectId][$role][$base] = []; } - if (!isset($this->subscriptions[$projectId][$role][$channel][$identifier])) { - $this->subscriptions[$projectId][$role][$channel][$identifier] = []; + if (!isset($this->subscriptions[$projectId][$role][$base][$identifier])) { + $this->subscriptions[$projectId][$role][$base][$identifier] = []; } - $this->subscriptions[$projectId][$role][$channel][$identifier][$subscriptionId] = $data; + + $channelData = $data; + $channelData['actions'] = $actions; + + $this->subscriptions[$projectId][$role][$base][$identifier][$subscriptionId] = $channelData; } } // Union channels/roles across all subscriptions on the connection; overwriting would // leave getSubscriptionMetadata and full unsubscribe operating on stale state. + // Channels are stored in *base* form here so they match subscription-tree keys. $existing = $this->connections[$identifier] ?? []; $existingChannels = $existing['channels'] ?? []; $existingRoles = $existing['roles'] ?? []; @@ -124,7 +172,7 @@ class Realtime extends MessagingAdapter 'projectId' => $projectId, 'roles' => \array_values(\array_unique(\array_merge($existingRoles, $roles))), 'userId' => $userId ?? ($existing['userId'] ?? ''), - 'channels' => \array_values(\array_unique(\array_merge($existingChannels, $channels))), + 'channels' => \array_values(\array_unique(\array_merge($existingChannels, $baseChannels))), ]; if (\array_key_exists('authorization', $existing)) { @@ -171,8 +219,16 @@ class Realtime extends MessagingAdapter 'queries' => $data['strings'] ?? [] ]; } - if (!\in_array($channel, $subscriptions[$subscriptionId]['channels'])) { - $subscriptions[$subscriptionId]['channels'][] = $channel; + + // Re-attach the action suffix so the original subscription channel + // (e.g. `documents.create`) is round-tripped on response paths and + // re-subscribe flows. `*` means no action was set — emit the base. + $actions = $data['actions'] ?? [self::ACTION_ALL]; + foreach ($actions as $action) { + $name = $action === self::ACTION_ALL ? $channel : $channel . '.' . $action; + if (!\in_array($name, $subscriptions[$subscriptionId]['channels'], true)) { + $subscriptions[$subscriptionId]['channels'][] = $name; + } } } } @@ -373,6 +429,7 @@ class Realtime extends MessagingAdapter } $payload = $event['data']['payload'] ?? []; + $events = $event['data']['events'] ?? []; foreach ($this->subscriptions[$event['project']] as $role => $subscriptionsByChannel) { foreach ($event['data']['channels'] as $channel) { @@ -389,6 +446,11 @@ class Realtime extends MessagingAdapter foreach ($subscriptions as $subscriptionId => $data) { $compiled = $data['compiled'] ?? ['type' => 'selectAll']; $strings = $data['strings'] ?? []; + $actions = $data['actions'] ?? [self::ACTION_ALL]; + + if (!self::matchesActions($actions, $events)) { + continue; + } if (RuntimeQuery::filter($compiled, $payload) !== null) { $matched[$subscriptionId] = $strings; @@ -408,6 +470,76 @@ class Realtime extends MessagingAdapter return $receivers; } + /** + * Tests whether any event in `$events` ends with one of the listed actions. + * + * Implements `containsAny` semantics (any-of match) over the events array, + * comparing only the trailing `.`-segment of each event so wildcard variants + * like `databases.*.collections.*.documents.*.create` match action `create`. + * `['*']` (or empty) means "no action filter" and short-circuits to true. + * + * @param array $actions Stored action filter for the subscription. + * @param array $events Event names from the published event. + * @return bool + */ + private static function matchesActions(array $actions, array $events): bool + { + if (empty($actions) || \in_array(self::ACTION_ALL, $actions, true)) { + return true; + } + + foreach ($events as $event) { + $lastDot = \strrpos($event, '.'); + if ($lastDot === false) { + continue; + } + if (\in_array(\substr($event, $lastDot + 1), $actions, true)) { + return true; + } + } + + return false; + } + + /** + * Splits a channel name into its base form and an action suffix. + * + * A trailing segment that matches one of {@see self::SUPPORTED_ACTIONS} is treated + * as an action filter and stripped from the channel; the remaining prefix becomes + * the base channel used for subscription-tree lookup. When no recognised suffix is + * present, the channel is returned unchanged with action {@see self::ACTION_ALL} + * (meaning "no action filter"). + * + * Examples: + * `documents.create` -> [`documents`, `create`] + * `databases.X.collections.Y.documents.Z.create` -> [`databases.X.collections.Y.documents.Z`, `create`] + * `documents` -> [`documents`, `*`] + * `account.create` -> already filtered out by convertChannels() + * + * @param string $channel + * @return array{0: string, 1: string} [baseChannel, action] + */ + public static function parseActionChannel(string $channel): array + { + $lastDot = \strrpos($channel, '.'); + if ($lastDot === false) { + return [$channel, self::ACTION_ALL]; + } + + $suffix = \substr($channel, $lastDot + 1); + if (!\in_array($suffix, self::SUPPORTED_ACTIONS, true)) { + return [$channel, self::ACTION_ALL]; + } + + $base = \substr($channel, 0, $lastDot); + if ($base === '') { + // Pathological — channel was just ".create"; leave it alone. + return [$channel, self::ACTION_ALL]; + } + + return [$base, $suffix]; + } + /** * Converts the channels from the Query Params into an array. * Also renames the account channel to account.USER_ID and removes all illegal account channel variations. diff --git a/tests/unit/Messaging/MessagingTest.php b/tests/unit/Messaging/MessagingTest.php index f48be46202..0706f8e89b 100644 --- a/tests/unit/Messaging/MessagingTest.php +++ b/tests/unit/Messaging/MessagingTest.php @@ -517,4 +517,264 @@ class MessagingTest extends TestCase $this->assertContains(Role::any()->toString(), $result['roles']); $this->assertContains(Role::team('123abc')->toString(), $result['roles']); } + + public function testParseActionChannel(): void + { + $this->assertSame(['documents', 'create'], Realtime::parseActionChannel('documents.create')); + $this->assertSame(['documents', 'update'], Realtime::parseActionChannel('documents.update')); + $this->assertSame(['documents', 'upsert'], Realtime::parseActionChannel('documents.upsert')); + $this->assertSame( + ['databases.X.collections.Y.documents.Z', 'create'], + Realtime::parseActionChannel('databases.X.collections.Y.documents.Z.create') + ); + + // No action suffix → unchanged with '*' default. + $this->assertSame(['documents', '*'], Realtime::parseActionChannel('documents')); + $this->assertSame(['documents.789', '*'], Realtime::parseActionChannel('documents.789')); + + // Unrecognised suffix (e.g. delete is not yet supported) → unchanged. + $this->assertSame(['documents.delete', '*'], Realtime::parseActionChannel('documents.delete')); + } + + public function testActionChannelFiltersByEventAction(): void + { + $realtime = new Realtime(); + + // Two subscriptions on the same connection: one filtered to creates only, + // one filtered to updates only. + $realtime->subscribe( + '1', + 1, + 'sub-create', + [Role::any()->toString()], + ['documents.create'], + ); + $realtime->subscribe( + '1', + 1, + 'sub-update', + [Role::any()->toString()], + ['documents.update'], + ); + + $createEvent = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['documents'], + 'events' => [ + 'databases.db.collections.col.documents.doc.create', + 'databases.*.collections.*.documents.*.create', + ], + 'payload' => ['$id' => 'doc'], + ], + ]; + + $updateEvent = $createEvent; + $updateEvent['data']['events'] = [ + 'databases.db.collections.col.documents.doc.update', + 'databases.*.collections.*.documents.*.update', + ]; + + // Create event should only deliver to sub-create. + $receivers = $realtime->getSubscribers($createEvent); + $this->assertCount(1, $receivers); + $this->assertArrayHasKey(1, $receivers); + $this->assertArrayHasKey('sub-create', $receivers[1]); + $this->assertArrayNotHasKey('sub-update', $receivers[1]); + + // Update event should only deliver to sub-update. + $receivers = $realtime->getSubscribers($updateEvent); + $this->assertCount(1, $receivers); + $this->assertArrayHasKey('sub-update', $receivers[1]); + $this->assertArrayNotHasKey('sub-create', $receivers[1]); + } + + public function testActionChannelHonorsResourceId(): void + { + $realtime = new Realtime(); + + // Subscribe to creates on a specific document only. + $realtime->subscribe( + '1', + 1, + 'sub-doc-create', + [Role::any()->toString()], + ['documents.789.create'], + ); + + // The base channel for `documents.789.create` is `documents.789`. + $event = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['documents.789'], + 'events' => [ + 'databases.db.collections.col.documents.789.create', + 'databases.*.collections.*.documents.*.create', + ], + 'payload' => ['$id' => '789'], + ], + ]; + + $receivers = $realtime->getSubscribers($event); + $this->assertCount(1, $receivers); + $this->assertArrayHasKey('sub-doc-create', $receivers[1]); + + // Update on the same document should not match. + $event['data']['events'] = [ + 'databases.db.collections.col.documents.789.update', + 'databases.*.collections.*.documents.*.update', + ]; + + $this->assertEmpty($realtime->getSubscribers($event)); + + // Create on a different document should not match (different base channel + // entirely; subscription tree key won't even line up). + $event['data']['channels'] = ['documents.999']; + $event['data']['events'] = [ + 'databases.db.collections.col.documents.999.create', + 'databases.*.collections.*.documents.*.create', + ]; + + $this->assertEmpty($realtime->getSubscribers($event)); + } + + public function testNonActionChannelStillReceivesAllEvents(): void + { + $realtime = new Realtime(); + + $realtime->subscribe( + '1', + 1, + 'sub-all', + [Role::any()->toString()], + ['documents'], + ); + + $event = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['documents'], + 'events' => [ + 'databases.db.collections.col.documents.doc.create', + ], + 'payload' => ['$id' => 'doc'], + ], + ]; + + $this->assertArrayHasKey(1, $realtime->getSubscribers($event)); + + $event['data']['events'] = ['databases.db.collections.col.documents.doc.update']; + $this->assertArrayHasKey(1, $realtime->getSubscribers($event)); + + $event['data']['events'] = ['databases.db.collections.col.documents.doc.upsert']; + $this->assertArrayHasKey(1, $realtime->getSubscribers($event)); + } + + public function testMixedActionAndBaseChannelInSameSubscription(): void + { + $realtime = new Realtime(); + + // Same sub-id covers `documents.create` (filtered) and `files` (unfiltered). + // After parsing they live under different base-channel keys with their own + // action metadata, so each gets its own filter behaviour. + $realtime->subscribe( + '1', + 1, + 'sub-mixed', + [Role::any()->toString()], + ['documents.create', 'files'], + ); + + // Create event on documents → matches. + $createDoc = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['documents'], + 'events' => ['databases.db.collections.col.documents.doc.create'], + 'payload' => [], + ], + ]; + $this->assertArrayHasKey(1, $realtime->getSubscribers($createDoc)); + + // Update event on documents → blocked by the action filter on the documents key. + $updateDoc = $createDoc; + $updateDoc['data']['events'] = ['databases.db.collections.col.documents.doc.update']; + $this->assertEmpty($realtime->getSubscribers($updateDoc)); + + // Files channel has no action filter — any action delivers. + $updateFile = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['files'], + 'events' => ['buckets.bucket.files.file.update'], + 'payload' => [], + ], + ]; + $this->assertArrayHasKey(1, $realtime->getSubscribers($updateFile)); + } + + public function testActionChannelMetadataRoundTrips(): void + { + $realtime = new Realtime(); + + $realtime->subscribe( + '1', + 1, + 'sub-create', + [Role::any()->toString()], + ['documents.create', 'files'], + ); + + $meta = $realtime->getSubscriptionMetadata(1); + + $this->assertArrayHasKey('sub-create', $meta); + $this->assertContains('documents.create', $meta['sub-create']['channels']); + $this->assertContains('files', $meta['sub-create']['channels']); + // Base form should NOT leak when an action was set. + $this->assertNotContains('documents', $meta['sub-create']['channels']); + } + + public function testMergingMultipleActionsOnSameBaseChannel(): void + { + $realtime = new Realtime(); + + // Subscribing to multiple actions on the same base merges their action lists + // onto a single tree entry. + $realtime->subscribe( + '1', + 1, + 'sub-multi', + [Role::any()->toString()], + ['documents.create', 'documents.update'], + ); + + $createEvent = [ + 'project' => '1', + 'roles' => [Role::any()->toString()], + 'data' => [ + 'channels' => ['documents'], + 'events' => ['databases.db.collections.col.documents.doc.create'], + 'payload' => [], + ], + ]; + $this->assertArrayHasKey(1, $realtime->getSubscribers($createEvent)); + + $updateEvent = $createEvent; + $updateEvent['data']['events'] = ['databases.db.collections.col.documents.doc.update']; + $this->assertArrayHasKey(1, $realtime->getSubscribers($updateEvent)); + + // Upsert should not match — neither create nor update covers it. + $upsertEvent = $createEvent; + $upsertEvent['data']['events'] = ['databases.db.collections.col.documents.doc.upsert']; + $this->assertEmpty($realtime->getSubscribers($upsertEvent)); + + $meta = $realtime->getSubscriptionMetadata(1); + $this->assertContains('documents.create', $meta['sub-multi']['channels']); + $this->assertContains('documents.update', $meta['sub-multi']['channels']); + } }