mirror of
https://github.com/appwrite/appwrite.git
synced 2026-05-26 13:51:13 +00:00
Enhance Realtime adapter with action channel support and tests
- Introduced ACTION_ALL and SUPPORTED_ACTIONS constants for better action handling. - Updated channel subscription logic to support action suffixes. - Added tests for action channel parsing and filtering in MessagingTest.
This commit is contained in:
@@ -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<string> $actions Stored action filter for the subscription.
|
||||
* @param array<string> $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.
|
||||
|
||||
@@ -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']);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user