Merge pull request #11202 from appwrite/realtime-query-subscriptions

This commit is contained in:
Jake Barnby
2026-02-03 11:44:13 +00:00
committed by GitHub
10 changed files with 1122 additions and 170 deletions
+1 -1
View File
@@ -4,7 +4,7 @@ namespace Appwrite\Messaging;
abstract class Adapter
{
abstract public function subscribe(string $projectId, mixed $identifier, array $roles, array $channels): void;
abstract public function subscribe(string $projectId, mixed $identifier, string $subscriptionId, array $roles, array $channels, array $queryGroup = []): void;
abstract public function unsubscribe(mixed $identifier): void;
abstract public function send(string $projectId, array $payload, array $events, array $channels, array $roles, array $options): void;
}
+201 -46
View File
@@ -29,13 +29,13 @@ class Realtime extends MessagingAdapter
*
* [PROJECT_ID] ->
* [ROLE_X] ->
* [CHANNEL_NAME_X] -> [CONNECTION_ID]
* [CHANNEL_NAME_Y] -> [CONNECTION_ID]
* [CHANNEL_NAME_Z] -> [CONNECTION_ID]
* [ROLE_Y] ->
* [CHANNEL_NAME_X] -> [CONNECTION_ID]
* [CHANNEL_NAME_Y] -> [CONNECTION_ID]
* [CHANNEL_NAME_Z] -> [CONNECTION_ID]
* [CHANNEL_NAME_X] ->
* [CONNECTION_ID] ->
* [SUB_ID] -> [query1, query2, ...] // Subscription with queries (AND logic)
*
* Each subscription ID maps to an array of query strings.
* Within a subscription: AND logic (all queries must match)
* Across subscriptions: OR logic (any subscription matching = send event)
*/
public array $subscriptions = [];
@@ -48,41 +48,108 @@ class Realtime extends MessagingAdapter
}
/**
* Adds a subscription.
* Adds a subscription with a specific subscription ID.
*
* @param string $projectId
* @param mixed $identifier
* @param array $roles
* @param array $channels
* @param array $queries
* @param mixed $identifier Connection ID
* @param string $subscriptionId Unique subscription ID
* @param array $roles User roles
* @param array $channels Channels to subscribe to (array of channel names)
* @param array $queryGroup Array of Query objects for this subscription (AND logic within subscription)
* @return void
*/
public function subscribe(string $projectId, mixed $identifier, array $roles, array $channels, array $queries = []): void
public function subscribe(string $projectId, mixed $identifier, string $subscriptionId, array $roles, array $channels, array $queryGroup = []): void
{
if (!isset($this->subscriptions[$projectId])) { // Init Project
$this->subscriptions[$projectId] = [];
}
foreach ($roles as $role) {
if (!isset($this->subscriptions[$projectId][$role])) { // Add user first connection
$this->subscriptions[$projectId][$role] = [];
}
foreach ($channels as $channel => $list) {
$this->subscriptions[$projectId][$role][$channel][$identifier] = true;
// Convert Query objects to strings for this subscription
$queryStrings = [];
if (empty($queryGroup)) {
// No queries means "listen to all events" - use select("*")
$queryStrings[] = Query::select(['*'])->toString();
} else {
foreach ($queryGroup as $query) {
/** @var Query $query */
$queryStrings[] = $query->toString();
}
}
foreach ($roles as $role) {
if (!isset($this->subscriptions[$projectId][$role])) {
$this->subscriptions[$projectId][$role] = [];
}
foreach ($channels as $channel) {
if (!isset($this->subscriptions[$projectId][$role][$channel])) {
$this->subscriptions[$projectId][$role][$channel] = [];
}
if (!isset($this->subscriptions[$projectId][$role][$channel][$identifier])) {
$this->subscriptions[$projectId][$role][$channel][$identifier] = [];
}
// Store subscription under subscription ID
$this->subscriptions[$projectId][$role][$channel][$identifier][$subscriptionId] = $queryStrings;
}
}
// Update connection info
$this->connections[$identifier] = [
'projectId' => $projectId,
'roles' => $roles,
'channels' => $channels,
'queries' => $queries
'channels' => $channels
];
}
/**
* Removes Subscription.
* Get subscription metadata for a connection.
* Retrieves subscription data including channels and queries directly from the subscriptions tree.
*
* @param mixed $connection Connection ID
* @return array Array of [subscriptionId => ['channels' => string[], 'queries' => string[]]]
*/
public function getSubscriptionMetadata(mixed $connection): array
{
$projectId = $this->connections[$connection]['projectId'] ?? null;
$roles = $this->connections[$connection]['roles'] ?? [];
$channels = $this->connections[$connection]['channels'] ?? [];
if (!$projectId || empty($roles) || empty($channels)) {
return [];
}
$subscriptions = [];
// Extract subscription data from subscriptions tree
foreach ($roles as $role) {
if (!isset($this->subscriptions[$projectId][$role])) {
continue;
}
foreach ($channels as $channel) {
if (!isset($this->subscriptions[$projectId][$role][$channel][$connection])) {
continue;
}
foreach ($this->subscriptions[$projectId][$role][$channel][$connection] as $subId => $queryStrings) {
if (!isset($subscriptions[$subId])) {
$subscriptions[$subId] = [
'channels' => [],
'queries' => $queryStrings
];
}
if (!in_array($channel, $subscriptions[$subId]['channels'])) {
$subscriptions[$subId]['channels'][] = $channel;
}
}
}
}
return $subscriptions;
}
/**
* Removes all subscriptions for a connection.
*
* @param mixed $connection
* @return void
@@ -91,10 +158,11 @@ class Realtime extends MessagingAdapter
{
$projectId = $this->connections[$connection]['projectId'] ?? '';
$roles = $this->connections[$connection]['roles'] ?? [];
$channels = $this->connections[$connection]['channels'] ?? [];
foreach ($roles as $role) {
foreach ($this->subscriptions[$projectId][$role] as $channel => $list) {
unset($this->subscriptions[$projectId][$role][$channel][$connection]); // Remove connection
foreach ($channels as $channel) {
unset($this->subscriptions[$projectId][$role][$channel][$connection]); // dropping connection will drop all subscriptions
if (empty($this->subscriptions[$projectId][$role][$channel])) {
unset($this->subscriptions[$projectId][$role][$channel]); // Remove channel when no connections
@@ -110,7 +178,9 @@ class Realtime extends MessagingAdapter
unset($this->subscriptions[$projectId]);
}
unset($this->connections[$connection]);
if (isset($this->connections[$connection])) {
unset($this->connections[$connection]);
}
}
/**
@@ -130,7 +200,8 @@ class Realtime extends MessagingAdapter
return array_key_exists($projectId, $this->subscriptions)
&& array_key_exists($role, $this->subscriptions[$projectId])
&& array_key_exists($channel, $this->subscriptions[$projectId][$role]);
&& array_key_exists($channel, $this->subscriptions[$projectId][$role])
&& !empty($this->subscriptions[$projectId][$role][$channel]);
}
/**
@@ -179,11 +250,10 @@ class Realtime extends MessagingAdapter
* - 1,121.328 ms (±0.84%) | 1,000,000 Connections / 10,000,000 Subscriptions
*
* @param array $event
* @return int[]|string[]
* @return array<int|string, array> Map of connection IDs to matched query groups
*/
public function getSubscribers(array $event): array
{
$receivers = [];
/**
* Check if project has subscriber.
@@ -207,17 +277,31 @@ class Realtime extends MessagingAdapter
/**
* Saving all connections that are allowed to receive this event.
*/
foreach (array_keys($this->subscriptions[$event['project']][$role][$channel]) as $id) {
/**
* To prevent duplicates, we save the connections as array keys.
*/
$queries = $this->connections[$id]['queries'] ?? [];
$payload = $event['data']['payload'] ?? [];
if (
empty($queries) ||
!empty(RuntimeQuery::filter($queries, $payload))
) {
$receivers[$id] = 0;
$payload = $event['data']['payload'] ?? [];
foreach ($this->subscriptions[$event['project']][$role][$channel] as $id => $subscriptions) {
$matchedSubscriptions = [];
// Process each subscription (OR logic across subscriptions)
foreach ($subscriptions as $subId => $queryStrings) {
$parsedQueries = [];
foreach ($queryStrings as $queryString) {
$parsed = Query::parseQueries([$queryString]);
$parsedQueries = array_merge($parsedQueries, $parsed);
}
// Check if this subscription matches (AND logic within subscription)
// Or if empty payload and select all as filter will return empty payload out of it even if it passed
$isEmptyPayloadAndSelectAll = RuntimeQuery::isSelectAll($parsedQueries[0]) && empty($payload);
if ($isEmptyPayloadAndSelectAll || !empty(RuntimeQuery::filter($parsedQueries, $payload))) {
$matchedSubscriptions[$subId] = $queryStrings;
}
}
// Only add connection to receivers if at least one subscription matched
if (!empty($matchedSubscriptions)) {
if (!isset($receivers[$id])) {
$receivers[$id] = [];
}
$receivers[$id] = array_merge($receivers[$id], $matchedSubscriptions);
}
}
break;
@@ -226,7 +310,7 @@ class Realtime extends MessagingAdapter
}
}
return array_keys($receivers);
return $receivers;
}
/**
@@ -258,17 +342,82 @@ class Realtime extends MessagingAdapter
}
/**
* Converts the queries from the Query Params into an array.
* @param array $queries
* @return array
* Constructs subscriptions from query parameters.
*
* Reconstructs subscription structure from query params where subscription indices can span multiple channels.
* Format: {channel}[subscriptionIndex][]=query1&{channel}[subscriptionIndex][]=query2
*
* Example:
* - tests[0][]=select(*) → subscription 0: channels=["tests"]
* - tests[1][]=equal(...) & prod[1][]=equal(...) → subscription 1: channels=["tests", "prod"]
*
* @param array $channelNames Array of channel names
* @param callable $getQueryParam Callable that takes a channel name and returns its query param value (null if not present)
* @return array Array indexed by subscription index: [index => ['channels' => string[], 'queries' => Query[]]]
* @throws QueryException
*/
public static function convertQueries(array $queries): array
public static function constructSubscriptions(array $channelNames, callable $getQueryParam): array
{
$subscriptionsByIndex = [];
foreach ($channelNames as $channel) {
$channelSubscriptions = $getQueryParam($channel);
// Backward compatibility: if no channel-specific query params, treat as subscription 0 with select("*")
if ($channelSubscriptions === null) {
if (!isset($subscriptionsByIndex[0])) {
$subscriptionsByIndex[0] = [
'channels' => [],
'queries' => []
];
}
$subscriptionsByIndex[0]['channels'][] = $channel;
if (empty($subscriptionsByIndex[0]['queries'])) {
$subscriptionsByIndex[0]['queries'] = [Query::select(['*'])];
}
continue;
}
if (!is_array($channelSubscriptions)) {
$channelSubscriptions = [$channelSubscriptions];
}
foreach ($channelSubscriptions as $subscriptionIndex => $subscription) {
if (!isset($subscriptionsByIndex[$subscriptionIndex])) {
$subscriptionsByIndex[$subscriptionIndex] = [
'channels' => [],
'queries' => []
];
}
if (!in_array($channel, $subscriptionsByIndex[$subscriptionIndex]['channels'])) {
$subscriptionsByIndex[$subscriptionIndex]['channels'][] = $channel;
}
if (empty($subscriptionsByIndex[$subscriptionIndex]['queries'])) {
$queriesToParse = is_array($subscription) ? $subscription : [$subscription];
$parsedQueries = self::convertQueries($queriesToParse);
$subscriptionsByIndex[$subscriptionIndex]['queries'] = $parsedQueries;
}
}
}
return $subscriptionsByIndex;
}
/**
* Converts the queries from the Query Params into an array.
* @param array|string $queries
* @return array
* @throws QueryException
*/
public static function convertQueries(mixed $queries): array
{
$queries = Query::parseQueries($queries);
$stack = $queries;
$allowedMethods = implode(', ', RuntimeQuery::ALLOWED_QUERIES);
while (!empty($stack)) {
/** `@var` Query $query */
/** @var Query $query */
$query = array_pop($stack);
$method = $query->getMethod();
if (!in_array($method, RuntimeQuery::ALLOWED_QUERIES, true)) {
@@ -277,6 +426,12 @@ class Realtime extends MessagingAdapter
"Query method '{$unsupportedMethod}' is not supported in Realtime queries. Allowed query methods are: {$allowedMethods}"
);
}
// Validate select queries - only select("*") is allowed
if ($method === Query::TYPE_SELECT) {
RuntimeQuery::validateSelectQuery($query);
}
if (in_array($method, [Query::TYPE_AND, Query::TYPE_OR], true)) {
$stack = array_merge($stack, $query->getValues());
}
+44 -1
View File
@@ -21,9 +21,44 @@ class RuntimeQuery extends Query
// Recursive checks
Query::TYPE_AND,
Query::TYPE_OR
Query::TYPE_OR,
// Special: select("*") means "listen to all events"
Query::TYPE_SELECT
];
/**
* Checks if a query is select("*") which means "listen to all events"
*
* @param Query $query
* @return bool
*/
public static function isSelectAll(Query $query): bool
{
return $query->getMethod() === Query::TYPE_SELECT
&& count($query->getValues()) === 1
&& $query->getValues()[0] === '*';
}
/**
* Validates a select query - only select("*") is allowed in Realtime
*
* @param Query $query
* @throws \InvalidArgumentException
*/
public static function validateSelectQuery(Query $query): void
{
if ($query->getMethod() !== Query::TYPE_SELECT) {
return;
}
if (!self::isSelectAll($query)) {
throw new \InvalidArgumentException(
'Only select("*") is allowed in Realtime queries. select("*") means "listen to all events".'
);
}
}
/**
* @param array<Query> $queries
* @param array<string, mixed> $payload
@@ -33,6 +68,14 @@ class RuntimeQuery extends Query
if (empty($queries)) {
return $payload;
}
// Check if select("*") is present - if so, return payload (match all)
foreach ($queries as $query) {
if (self::isSelectAll($query)) {
return $payload;
}
}
// multiple queries follows and condition
foreach ($queries as $query) {
if (!self::evaluateFilter($query, $payload)) {