diff --git a/app/realtime.php b/app/realtime.php index 1f0c4300a8..97e7465683 100644 --- a/app/realtime.php +++ b/app/realtime.php @@ -431,15 +431,20 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats, ] ]; - $subscribers = $realtime->getSubscribers($event); // [connectionId => [subId => queries]] + $subscribers = $realtime->getSubscribers($event); - // For test events, send to all connections with their matched subscription queries - foreach ($subscribers as $connectionId => $matchedSubscriptions) { + $groups = []; + foreach ($subscribers as $id => $matched) { + $key = implode(',', array_keys($matched)); + $groups[$key]['ids'][] = $id; + $groups[$key]['subscriptions'] = array_keys($matched); + } + + foreach ($groups as $group) { $data = $event['data']; - // Send matched subscription IDs - $data['subscriptions'] = array_keys($matchedSubscriptions); + $data['subscriptions'] = $group['subscriptions']; - $server->send([$connectionId], json_encode([ + $server->send($group['ids'], json_encode([ 'type' => 'event', 'data' => $data ])); @@ -484,18 +489,18 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats, $roles = $user->getRoles($database->getAuthorization()); $authorization = $realtime->connections[$connection]['authorization'] ?? null; - $subscriptionMetadata = $realtime->getSubscriptionMetadata($connection); + $meta = $realtime->getSubscriptionMetadata($connection); $realtime->unsubscribe($connection); - foreach ($subscriptionMetadata as $subscriptionId => $metadata) { - $queries = Query::parseQueries($metadata['queries'] ?? []); + foreach ($meta as $subscriptionId => $subscription) { + $queries = Query::parseQueries($subscription['queries'] ?? []); $realtime->subscribe( $projectId, $connection, $subscriptionId, $roles, - $metadata['channels'] ?? [], + $subscription['channels'] ?? [], $queries ); } @@ -507,35 +512,38 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats, } } - $receivers = $realtime->getSubscribers($event); // [connectionId => [subId => queries]] + $receivers = $realtime->getSubscribers($event); if (Http::isDevelopment() && !empty($receivers)) { Console::log("[Debug][Worker {$workerId}] Receivers: " . count($receivers)); - Console::log("[Debug][Worker {$workerId}] Receivers Connection IDs: " . json_encode(array_keys($receivers))); - Console::log("[Debug][Worker {$workerId}] Event Query: " . json_encode(array_values($receivers))); + Console::log("[Debug][Worker {$workerId}] Connection IDs: " . json_encode(array_keys($receivers))); + Console::log("[Debug][Worker {$workerId}] Matched: " . json_encode(array_values($receivers))); Console::log("[Debug][Worker {$workerId}] Event: " . $payload); } - $totalMessages = 0; - - foreach ($receivers as $connectionId => $matchedSubscriptions) { - $data = $event['data']; - // Send matched subscription IDs - $data['subscriptions'] = array_keys($matchedSubscriptions); - - $server->send( - [$connectionId], - json_encode([ - 'type' => 'event', - 'data' => $data - ]) - ); - $totalMessages++; + // Group connections by matched subscription IDs for batch sending + $groups = []; + foreach ($receivers as $id => $matched) { + $key = implode(',', array_keys($matched)); + $groups[$key]['ids'][] = $id; + $groups[$key]['subscriptions'] = array_keys($matched); } - if ($totalMessages > 0) { - $register->get('telemetry.messageSentCounter')->add($totalMessages); - $stats->incr($event['project'], 'messages', $totalMessages); + $total = 0; + foreach ($groups as $group) { + $data = $event['data']; + $data['subscriptions'] = $group['subscriptions']; + + $server->send($group['ids'], json_encode([ + 'type' => 'event', + 'data' => $data + ])); + $total += count($group['ids']); + } + + if ($total > 0) { + $register->get('telemetry.messageSentCounter')->add($total); + $stats->incr($event['project'], 'messages', $total); } }); } catch (Throwable $th) { @@ -624,21 +632,19 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, throw new Exception(Exception::REALTIME_POLICY_VIOLATION, 'Missing channels'); } - // Reconstruct subscriptions from query params using helper method - $channelNames = array_keys($channels); + $names = array_keys($channels); try { - $subscriptionsByIndex = Realtime::constructSubscriptions( - $channelNames, + $subscriptions = Realtime::constructSubscriptions( + $names, fn ($channel) => $request->getQuery($channel, null) ); } catch (QueryException $e) { throw new Exception(Exception::REALTIME_POLICY_VIOLATION, $e->getMessage()); } - // Generate subscription IDs and subscribe - $subscriptionMapping = []; - foreach ($subscriptionsByIndex as $index => $subscription) { + $mapping = []; + foreach ($subscriptions as $index => $subscription) { $subscriptionId = ID::unique(); $realtime->subscribe( @@ -647,10 +653,10 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $subscriptionId, $roles, $subscription['channels'], - $subscription['queries'] // Query objects + $subscription['queries'] ); - $subscriptionMapping[$index] = $subscriptionId; + $mapping[$index] = $subscriptionId; } $realtime->connections[$connection]['authorization'] = $authorization; @@ -660,8 +666,8 @@ $server->onOpen(function (int $connection, SwooleRequest $request) use ($server, $server->send([$connection], json_encode([ 'type' => 'connected', 'data' => [ - 'channels' => $channelNames, - 'subscriptions' => $subscriptionMapping, + 'channels' => $names, + 'subscriptions' => $mapping, 'user' => $user ] ])); @@ -791,32 +797,29 @@ $server->onMessage(function (int $connection, string $message) use ($server, $re } $roles = $user->getRoles($database->getAuthorization()); - $channelNames = $realtime->connections[$connection]['channels'] ?? []; - $channels = Realtime::convertChannels(array_flip($channelNames), $user->getId()); $authorization = $realtime->connections[$connection]['authorization'] ?? null; $projectId = $realtime->connections[$connection]['projectId'] ?? null; - $subscriptionMetadata = $realtime->getSubscriptionMetadata($connection); + $meta = $realtime->getSubscriptionMetadata($connection); $realtime->unsubscribe($connection); if (!empty($projectId)) { - foreach ($subscriptionMetadata as $subscriptionId => $metadata) { - $queries = Query::parseQueries($metadata['queries'] ?? []); + foreach ($meta as $subscriptionId => $subscription) { + $queries = Query::parseQueries($subscription['queries'] ?? []); $realtime->subscribe( $projectId, $connection, $subscriptionId, $roles, - $metadata['channels'] ?? [], + $subscription['channels'] ?? [], $queries ); } } - // Restore authorization after subscribe if ($authorization !== null) { $realtime->connections[$connection]['authorization'] = $authorization; } diff --git a/src/Appwrite/Messaging/Adapter/Realtime.php b/src/Appwrite/Messaging/Adapter/Realtime.php index 2bd0bc9cb8..51aad2a635 100644 --- a/src/Appwrite/Messaging/Adapter/Realtime.php +++ b/src/Appwrite/Messaging/Adapter/Realtime.php @@ -31,9 +31,9 @@ class Realtime extends MessagingAdapter * [ROLE_X] -> * [CHANNEL_NAME_X] -> * [CONNECTION_ID] -> - * [SUB_ID] -> ['strings' => [...], 'parsed' => [...]] + * [SUB_ID] -> ['strings' => [...], 'compiled' => [...]] * - * Each subscription ID maps to query strings (for metadata) and pre-parsed Query objects (for filtering). + * 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) */ @@ -73,25 +73,18 @@ class Realtime extends MessagingAdapter $this->subscriptions[$projectId] = []; } - // Convert Query objects to strings and store both for this subscription - $queryStrings = []; - $parsedQueries = []; + $strings = []; if (empty($queryGroup)) { - // No queries means "listen to all events" - use select("*") - $selectAll = Query::select(['*']); - $queryStrings[] = $selectAll->toString(); - $parsedQueries[] = $selectAll; + $strings[] = Query::select(['*'])->toString(); } else { foreach ($queryGroup as $query) { - /** @var Query $query */ - $queryStrings[] = $query->toString(); - $parsedQueries[] = $query; + $strings[] = $query->toString(); } } - $subscriptionData = [ - 'strings' => $queryStrings, - 'parsed' => $parsedQueries, + $data = [ + 'strings' => $strings, + 'compiled' => RuntimeQuery::compile($queryGroup), ]; foreach ($roles as $role) { @@ -106,7 +99,7 @@ class Realtime extends MessagingAdapter if (!isset($this->subscriptions[$projectId][$role][$channel][$identifier])) { $this->subscriptions[$projectId][$role][$channel][$identifier] = []; } - $this->subscriptions[$projectId][$role][$channel][$identifier][$subscriptionId] = $subscriptionData; + $this->subscriptions[$projectId][$role][$channel][$identifier][$subscriptionId] = $data; } } @@ -148,15 +141,15 @@ class Realtime extends MessagingAdapter continue; } - foreach ($this->subscriptions[$projectId][$role][$channel][$connection] as $subId => $subscriptionData) { - if (!isset($subscriptions[$subId])) { - $subscriptions[$subId] = [ + foreach ($this->subscriptions[$projectId][$role][$channel][$connection] as $subscriptionId => $data) { + if (!isset($subscriptions[$subscriptionId])) { + $subscriptions[$subscriptionId] = [ 'channels' => [], - 'queries' => $subscriptionData['strings'] ?? [] + 'queries' => $data['strings'] ?? [] ]; } - if (!\in_array($channel, $subscriptions[$subId]['channels'])) { - $subscriptions[$subId]['channels'][] = $channel; + if (!\in_array($channel, $subscriptions[$subscriptionId]['channels'])) { + $subscriptions[$subscriptionId]['channels'][] = $channel; } } } @@ -259,12 +252,10 @@ class Realtime extends MessagingAdapter * Identifies the receivers of all subscriptions, based on the permissions and event. * * Example of performance with an event with user:XXX permissions and with X users spread across 10 different channels: - * - 0.014 ms (±6.88%) | 10 Connections / 100 Subscriptions - * - 0.070 ms (±3.71%) | 100 Connections / 1,000 Subscriptions - * - 0.846 ms (±2.74%) | 1,000 Connections / 10,000 Subscriptions - * - 10.866 ms (±1.01%) | 10,000 Connections / 100,000 Subscriptions - * - 110.201 ms (±2.32%) | 100,000 Connections / 1,000,000 Subscriptions - * - 1,121.328 ms (±0.84%) | 1,000,000 Connections / 10,000,000 Subscriptions + * - 0.013 ms | 10 Connections / 100 Subscriptions + * - 0.14 ms | 100 Connections / 1,000 Subscriptions + * - 1.5 ms | 1,000 Connections / 10,000 Subscriptions + * - 15 ms | 10,000 Connections / 100,000 Subscriptions * * @param array $event * @return array Map of connection IDs to matched query groups @@ -272,55 +263,39 @@ class Realtime extends MessagingAdapter public function getSubscribers(array $event): array { $receivers = []; - /** - * Check if project has subscriber. - */ - if (isset($this->subscriptions[$event['project']])) { - /** - * Iterate through each role. - */ - foreach ($this->subscriptions[$event['project']] as $role => $subscription) { - /** - * Iterate through each channel. - */ - foreach ($event['data']['channels'] as $channel) { - /** - * Check if channel has subscriber. Also taking care of the role in the event and the wildcard role. - */ - if ( - \array_key_exists($channel, $this->subscriptions[$event['project']][$role]) - && (\in_array($role, $event['roles']) || \in_array(Role::any()->toString(), $event['roles'])) - ) { - /** - * Saving all connections that are allowed to receive this event. - */ - $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 => $subscriptionData) { - // Use pre-parsed queries instead of re-parsing on every event - $parsedQueries = $subscriptionData['parsed'] ?? []; - $queryStrings = $subscriptionData['strings'] ?? []; + if (!isset($this->subscriptions[$event['project']])) { + return $receivers; + } - // 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 = !empty($parsedQueries) && RuntimeQuery::isSelectAll($parsedQueries[0]) && empty($payload); - if ($isEmptyPayloadAndSelectAll || !empty(RuntimeQuery::filter($parsedQueries, $payload))) { - $matchedSubscriptions[$subId] = $queryStrings; - } - } + $payload = $event['data']['payload'] ?? []; - // Only add connection to receivers if at least one subscription matched - if (!empty($matchedSubscriptions)) { - if (!isset($receivers[$id])) { - $receivers[$id] = []; - } - $receivers[$id] += $matchedSubscriptions; - } + foreach ($this->subscriptions[$event['project']] as $role => $subscriptionsByChannel) { + foreach ($event['data']['channels'] as $channel) { + if ( + !\array_key_exists($channel, $subscriptionsByChannel) + || (!\in_array($role, $event['roles']) && !\in_array(Role::any()->toString(), $event['roles'])) + ) { + continue; + } + + foreach ($subscriptionsByChannel[$channel] as $id => $subscriptions) { + $matched = []; + + foreach ($subscriptions as $subscriptionId => $data) { + $compiled = $data['compiled'] ?? ['type' => 'selectAll']; + $strings = $data['strings'] ?? []; + + if (RuntimeQuery::filter($compiled, $payload) !== null) { + $matched[$subscriptionId] = $strings; } - break; + } + + if (!empty($matched)) { + if (!isset($receivers[$id])) { + $receivers[$id] = []; + } + $receivers[$id] += $matched; } } } @@ -360,65 +335,50 @@ class Realtime extends MessagingAdapter /** * 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[]]] + * @param array $channelNames + * @param callable $getQueryParam + * @return array [index => ['channels' => string[], 'queries' => Query[]]] * @throws QueryException */ public static function constructSubscriptions(array $channelNames, callable $getQueryParam): array { - $subscriptionsByIndex = []; + $subscriptions = []; foreach ($channelNames as $channel) { - $channelSubscriptions = $getQueryParam(str_replace(".", "_", $channel)); + $params = $getQueryParam(\str_replace('.', '_', $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' => [] - ]; + if ($params === null) { + if (!isset($subscriptions[0])) { + $subscriptions[0] = ['channels' => [], 'queries' => []]; } - $subscriptionsByIndex[0]['channels'][] = $channel; - if (empty($subscriptionsByIndex[0]['queries'])) { - $subscriptionsByIndex[0]['queries'] = [Query::select(['*'])]; + $subscriptions[0]['channels'][] = $channel; + if (empty($subscriptions[0]['queries'])) { + $subscriptions[0]['queries'] = [Query::select(['*'])]; } continue; } - if (!is_array($channelSubscriptions)) { - $channelSubscriptions = [$channelSubscriptions]; + if (!is_array($params)) { + $params = [$params]; } - foreach ($channelSubscriptions as $subscriptionIndex => $subscription) { - if (!isset($subscriptionsByIndex[$subscriptionIndex])) { - $subscriptionsByIndex[$subscriptionIndex] = [ - 'channels' => [], - 'queries' => [] - ]; + foreach ($params as $index => $slot) { + if (!isset($subscriptions[$index])) { + $subscriptions[$index] = ['channels' => [], 'queries' => []]; } - if (!in_array($channel, $subscriptionsByIndex[$subscriptionIndex]['channels'])) { - $subscriptionsByIndex[$subscriptionIndex]['channels'][] = $channel; + if (!in_array($channel, $subscriptions[$index]['channels'])) { + $subscriptions[$index]['channels'][] = $channel; } - if (empty($subscriptionsByIndex[$subscriptionIndex]['queries'])) { - $queriesToParse = is_array($subscription) ? $subscription : [$subscription]; - $parsedQueries = self::convertQueries($queriesToParse); - $subscriptionsByIndex[$subscriptionIndex]['queries'] = $parsedQueries; + if (empty($subscriptions[$index]['queries'])) { + $raw = is_array($slot) ? $slot : [$slot]; + $subscriptions[$index]['queries'] = self::convertQueries($raw); } } } - return $subscriptionsByIndex; + return $subscriptions; } /** @@ -431,19 +391,18 @@ class Realtime extends MessagingAdapter { $queries = Query::parseQueries($queries); $stack = $queries; - $allowedMethods = implode(', ', RuntimeQuery::ALLOWED_QUERIES); + $allowed = implode(', ', RuntimeQuery::ALLOWED_QUERIES); + while (!empty($stack)) { - /** @var Query $query */ $query = array_pop($stack); $method = $query->getMethod(); + if (!in_array($method, RuntimeQuery::ALLOWED_QUERIES, true)) { - $unsupportedMethod = $method; throw new QueryException( - "Query method '{$unsupportedMethod}' is not supported in Realtime queries. Allowed query methods are: {$allowedMethods}" + "Query method '{$method}' is not supported in Realtime queries. Allowed: {$allowed}" ); } - // Validate select queries - only select("*") is allowed if ($method === Query::TYPE_SELECT) { RuntimeQuery::validateSelectQuery($query); } diff --git a/src/Appwrite/Utopia/Database/RuntimeQuery.php b/src/Appwrite/Utopia/Database/RuntimeQuery.php index f959e9b573..369006d9df 100644 --- a/src/Appwrite/Utopia/Database/RuntimeQuery.php +++ b/src/Appwrite/Utopia/Database/RuntimeQuery.php @@ -4,6 +4,11 @@ namespace Appwrite\Utopia\Database; use Utopia\Database\Query; +/** + * RuntimeQuery handles real-time query filtering for Appwrite's Realtime subscriptions. + * + * Queries are pre-compiled at subscription time for fast evaluation during message delivery. + */ class RuntimeQuery extends Query { public const ALLOWED_QUERIES = [ @@ -28,16 +33,19 @@ class RuntimeQuery extends Query ]; /** - * Checks if a query is select("*") which means "listen to all events" + * Checks if a query is a select("*") query. * * @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] === '*'; + if ($query->getMethod() !== Query::TYPE_SELECT) { + return false; + } + + $values = $query->getValues(); + return count($values) === 1 && $values[0] === '*'; } /** @@ -60,105 +68,222 @@ class RuntimeQuery extends Query } /** + * Pre-compile queries into an optimized format for fast evaluation. + * Call this once when subscription is created, store the result. + * * @param array $queries - * @param array $payload + * @return array Compiled query structure with 'type' key */ - public static function filter(array $queries, array $payload): array + public static function compile(array $queries): array { if (empty($queries)) { - return $payload; + return ['type' => 'selectAll']; } - // Check if select("*") is present - if so, return payload (match all) + // Check for select("*") upfront foreach ($queries as $query) { - if (self::isSelectAll($query)) { - return $payload; + if ($query->getMethod() === Query::TYPE_SELECT) { + $values = $query->getValues(); + if (count($values) === 1 && $values[0] === '*') { + return ['type' => 'selectAll']; + } } } - // multiple queries follows and condition + // Compile queries into flat structure + $compiled = [ + 'type' => 'filter', + 'conditions' => [], + 'attributes' => [], + 'hasOr' => false, + ]; + foreach ($queries as $query) { - if (!self::evaluateFilter($query, $payload)) { - return []; - }; + $condition = self::compileCondition($query); + $compiled['conditions'][] = $condition; + self::extractAttributes($condition, $compiled['attributes'], $compiled['hasOr']); } + + $compiled['attributes'] = array_unique($compiled['attributes']); + + return $compiled; + } + + /** + * Compile a single query condition into an optimized array format. + */ + private static function compileCondition(Query $query): array + { + $method = $query->getMethod(); + + if ($method === Query::TYPE_AND) { + return [ + 'op' => 'AND', + 'conditions' => array_map([self::class, 'compileCondition'], $query->getValues()), + ]; + } + + if ($method === Query::TYPE_OR) { + return [ + 'op' => 'OR', + 'conditions' => array_map([self::class, 'compileCondition'], $query->getValues()), + ]; + } + + return [ + 'op' => $method, + 'attr' => $query->getAttribute(), + 'values' => $query->getValues(), + ]; + } + + /** + * Extract all attribute names from a compiled condition tree. + * Also tracks whether any OR conditions exist. + */ + private static function extractAttributes(array $condition, array &$attributes, bool &$hasOr): void + { + if (isset($condition['op']) && $condition['op'] === 'OR') { + $hasOr = true; + } + if (isset($condition['attr'])) { + $attributes[] = $condition['attr']; + } + if (isset($condition['conditions'])) { + foreach ($condition['conditions'] as $sub) { + self::extractAttributes($sub, $attributes, $hasOr); + } + } + } + + /** + * Fast filter using pre-compiled query structure. + * + * @param array $compiled Result from compile() + * @param array $payload Event payload + * @return array|null Null if no match, payload if match + */ + public static function filter(array $compiled, array $payload): ?array + { + // Fast path for select("*") subscriptions + if ($compiled['type'] === 'selectAll') { + return $payload; + } + + // Quick rejection: if payload is missing any required attribute, fail fast + // Skip this optimization when OR conditions exist (OR can match with partial attributes) + if (empty($compiled['hasOr'])) { + foreach ($compiled['attributes'] as $attr) { + if (!isset($payload[$attr]) && !\array_key_exists($attr, $payload)) { + return null; + } + } + } + + // Evaluate all conditions (AND logic at top level) + foreach ($compiled['conditions'] as $condition) { + if (!self::evaluateCondition($condition, $payload)) { + return null; + } + } + return $payload; } - private static function evaluateFilter(Query $query, array $payload): bool + /** + * Evaluate a single compiled condition against a payload. + */ + private static function evaluateCondition(array $condition, array $payload): bool { - $attribute = $query->getAttribute(); - $method = $query->getMethod(); - $values = $query->getValues(); + $op = $condition['op']; - // during 'and' and 'or' attribute will not be present - switch ($method) { - case Query::TYPE_AND: - // All subqueries must evaluate to true - foreach ($query->getValues() as $subquery) { - if (!self::evaluateFilter($subquery, $payload)) { + // Handle AND/OR + if ($op === 'AND') { + foreach ($condition['conditions'] as $sub) { + if (!self::evaluateCondition($sub, $payload)) { + return false; + } + } + return true; + } + + if ($op === 'OR') { + foreach ($condition['conditions'] as $sub) { + if (self::evaluateCondition($sub, $payload)) { + return true; + } + } + return false; + } + + // Leaf condition - direct comparison + $attr = $condition['attr']; + + if (!\array_key_exists($attr, $payload)) { + return false; + } + + $value = $payload[$attr]; + $targets = $condition['values']; + + // Inlined comparisons - no closures, no method calls + switch ($op) { + case Query::TYPE_EQUAL: + foreach ($targets as $target) { + if ($value === $target) { + return true; + } + } + return false; + + case Query::TYPE_NOT_EQUAL: + foreach ($targets as $target) { + if ($value === $target) { return false; } } return true; - case Query::TYPE_OR: - // At least one subquery must evaluate to true - foreach ($query->getValues() as $subquery) { - if (self::evaluateFilter($subquery, $payload)) { + case Query::TYPE_LESSER: + foreach ($targets as $target) { + if ($value < $target) { return true; } } return false; - } - - $hasAttribute = \array_key_exists($attribute, $payload); - if (!$hasAttribute) { - return false; - } - - // null can be a value as well - $payloadAttributeValue = $payload[$attribute]; - switch ($method) { - case Query::TYPE_EQUAL: - return self::anyMatch($values, fn ($value) => $payloadAttributeValue === $value); - - case Query::TYPE_NOT_EQUAL: - return !self::anyMatch($values, fn ($value) => $payloadAttributeValue === $value); - - case Query::TYPE_LESSER: - return self::anyMatch($values, fn ($value) => $payloadAttributeValue < $value); case Query::TYPE_LESSER_EQUAL: - return self::anyMatch($values, fn ($value) => $payloadAttributeValue <= $value); + foreach ($targets as $target) { + if ($value <= $target) { + return true; + } + } + return false; case Query::TYPE_GREATER: - return self::anyMatch($values, fn ($value) => $payloadAttributeValue > $value); + foreach ($targets as $target) { + if ($value > $target) { + return true; + } + } + return false; case Query::TYPE_GREATER_EQUAL: - return self::anyMatch($values, fn ($value) => $payloadAttributeValue >= $value); + foreach ($targets as $target) { + if ($value >= $target) { + return true; + } + } + return false; - // attribute must be present and should be explicitly null case Query::TYPE_IS_NULL: - return $payloadAttributeValue === null; + return $value === null; case Query::TYPE_IS_NOT_NULL: - return $payloadAttributeValue !== null; + return $value !== null; default: - throw new \InvalidArgumentException( - "Unsupported query method: {$method}" - ); + return false; } } - - private static function anyMatch(array $values, callable $fn): bool - { - foreach ($values as $value) { - if ($fn($value)) { - return true; - } - } - return false; - } } diff --git a/tests/unit/Utopia/Database/Query/RuntimeQueryTest.php b/tests/unit/Utopia/Database/Query/RuntimeQueryTest.php index 51d3a307da..f7d73eb287 100644 --- a/tests/unit/Utopia/Database/Query/RuntimeQueryTest.php +++ b/tests/unit/Utopia/Database/Query/RuntimeQueryTest.php @@ -16,10 +16,19 @@ class RuntimeQueryTest extends TestCase { } + /** + * Helper to compile and filter queries in one step for tests. + */ + private function compileAndFilter(array $queries, array $payload): ?array + { + $compiled = RuntimeQuery::compile($queries); + return RuntimeQuery::filter($compiled, $payload); + } + public function testFilterEmptyQueries(): void { $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([], $payload); + $result = $this->compileAndFilter([], $payload); $this->assertEquals($payload, $result); } @@ -27,15 +36,15 @@ class RuntimeQueryTest extends TestCase { $queries = [Query::equal('name', ['Jane'])]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter($queries, $payload); + $this->assertNull($result); } public function testFilterWithMatchingQuery(): void { $queries = [Query::equal('name', ['John'])]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); + $result = $this->compileAndFilter($queries, $payload); $this->assertEquals($payload, $result); } @@ -44,7 +53,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('name', ['John']); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -52,15 +61,15 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('name', ['Jane']); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testEqualMultipleValuesMatch(): void { $query = Query::equal('status', ['active', 'pending', 'approved']); $payload = ['status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -68,15 +77,15 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('status', ['active', 'pending', 'approved']); $payload = ['status' => 'rejected']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testEqualNumericValues(): void { $query = Query::equal('age', [30, 25, 35]); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -84,7 +93,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('active', [true]); $payload = ['active' => true]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -92,8 +101,8 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('missing', ['value']); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } // TYPE_NOT_EQUAL tests @@ -101,7 +110,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::notEqual('name', ['Jane']); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -109,8 +118,8 @@ class RuntimeQueryTest extends TestCase { $query = Query::notEqual('name', ['John']); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testNotEqualMultipleValues(): void @@ -119,13 +128,13 @@ class RuntimeQueryTest extends TestCase // and Query::parse will be done first and parse doesn't allow multiple notEqual values $query = Query::notEqual('status', ['rejected', 'cancelled']); $payload = ['status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); $query = Query::notEqual('status', ['active', 'pending']); $payload = ['status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } // TYPE_LESSER tests @@ -133,7 +142,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThan('age', 30); $payload = ['age' => 25]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -141,16 +150,16 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThan('age', 30); $payload = ['age' => 35]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testLesserEqualValue(): void { $query = Query::lessThan('age', 30); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testLesserMultipleValues(): void @@ -159,7 +168,7 @@ class RuntimeQueryTest extends TestCase // This test uses a single value as Query class requires $query = Query::lessThan('age', 30); $payload = ['age' => 25]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -167,7 +176,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThan('name', 'M'); $payload = ['name' => 'A']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -176,7 +185,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThanEqual('age', 30); $payload = ['age' => 25]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -184,7 +193,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThanEqual('age', 30); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -192,8 +201,8 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThanEqual('age', 30); $payload = ['age' => 35]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testLesserEqualMultipleValues(): void @@ -201,7 +210,7 @@ class RuntimeQueryTest extends TestCase // Note: Query::lessThanEqual only accepts single value $query = Query::lessThanEqual('age', 30); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -210,7 +219,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThan('age', 30); $payload = ['age' => 35]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -218,16 +227,16 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThan('age', 30); $payload = ['age' => 25]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testGreaterEqualValue(): void { $query = Query::greaterThan('age', 30); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testGreaterMultipleValues(): void @@ -235,7 +244,7 @@ class RuntimeQueryTest extends TestCase // Note: Query::greaterThan only accepts single value $query = Query::greaterThan('age', 20); $payload = ['age' => 35]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -244,7 +253,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThanEqual('age', 30); $payload = ['age' => 35]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -252,7 +261,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThanEqual('age', 30); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -260,8 +269,8 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThanEqual('age', 30); $payload = ['age' => 25]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testGreaterEqualMultipleValues(): void @@ -269,7 +278,7 @@ class RuntimeQueryTest extends TestCase // Note: Query::greaterThanEqual only accepts single value $query = Query::greaterThanEqual('age', 20); $payload = ['age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -278,7 +287,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::isNull('description'); $payload = ['description' => null]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -286,16 +295,16 @@ class RuntimeQueryTest extends TestCase { $query = Query::isNull('description'); $payload = ['description' => 'Some text']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testIsNullMissingAttribute(): void { $query = Query::isNull('missing'); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } // TYPE_IS_NOT_NULL tests @@ -303,7 +312,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::isNotNull('description'); $payload = ['description' => 'Some text']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -311,16 +320,16 @@ class RuntimeQueryTest extends TestCase { $query = Query::isNotNull('description'); $payload = ['description' => null]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testIsNotNullMissingAttribute(): void { $query = Query::isNotNull('missing'); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } // TYPE_AND tests @@ -331,7 +340,7 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [30]) ]); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -342,8 +351,8 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [25]) ]); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testAndAllFail(): void @@ -353,8 +362,8 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [25]) ]); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testAndMultipleConditions(): void @@ -365,7 +374,7 @@ class RuntimeQueryTest extends TestCase Query::isNotNull('email') ]); $payload = ['status' => 'active', 'age' => 25, 'email' => 'test@example.com']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -379,7 +388,7 @@ class RuntimeQueryTest extends TestCase ]) ]); $payload = ['name' => 'John', 'age' => 30, 'status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -391,7 +400,7 @@ class RuntimeQueryTest extends TestCase Query::equal('name', ['Jane']) ]); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -402,7 +411,7 @@ class RuntimeQueryTest extends TestCase Query::equal('status', ['pending']) ]); $payload = ['status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -413,8 +422,8 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [25]) ]); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testOrMultipleConditions(): void @@ -425,7 +434,7 @@ class RuntimeQueryTest extends TestCase Query::equal('status', ['approved']) ]); $payload = ['status' => 'pending']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -439,7 +448,7 @@ class RuntimeQueryTest extends TestCase ]) ]); $payload = ['name' => 'Bob']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -450,10 +459,36 @@ class RuntimeQueryTest extends TestCase Query::equal('email', ['john@example.com']) ]); $payload = ['name' => 'Jane', 'email' => 'john@example.com']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } + public function testOrWithMissingAttributeInOneBranch(): void + { + // OR should match when one branch's attribute is missing but another branch matches + $query = Query::or([ + Query::equal('name', ['John']), + Query::equal('email', ['john@example.com']) + ]); + // Payload only has email, not name - should still match via email branch + $payload = ['email' => 'john@example.com']; + $result = $this->compileAndFilter([$query], $payload); + $this->assertEquals($payload, $result); + } + + public function testOrWithMissingAttributeNoMatch(): void + { + // OR should not match when the only matching branch has missing attribute + $query = Query::or([ + Query::equal('name', ['John']), + Query::equal('email', ['john@example.com']) + ]); + // Payload only has name but it doesn't match - should not match + $payload = ['name' => 'Jane']; + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); + } + // Complex combinations public function testAndOrCombination(): void { @@ -465,7 +500,7 @@ class RuntimeQueryTest extends TestCase ]) ]); $payload = ['type' => 'user', 'status' => 'active']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -482,7 +517,7 @@ class RuntimeQueryTest extends TestCase ]) ]); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -494,7 +529,7 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [30]) ]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); + $result = $this->compileAndFilter($queries, $payload); $this->assertEquals($payload, $result); } @@ -505,9 +540,9 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [25]) ]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); + $result = $this->compileAndFilter($queries, $payload); // With AND logic, if first matches but second doesn't, should return empty - $this->assertEquals([], $result); + $this->assertNull($result); } public function testMultipleQueriesSecondMatches(): void @@ -517,9 +552,9 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [30]) ]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); + $result = $this->compileAndFilter($queries, $payload); // With AND logic, if second matches but first doesn't, should return empty - $this->assertEquals([], $result); + $this->assertNull($result); } public function testMultipleQueriesNoneMatch(): void @@ -529,23 +564,23 @@ class RuntimeQueryTest extends TestCase Query::equal('age', [25]) ]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter($queries, $payload); + $this->assertNull($result); } public function testEmptyPayload(): void { $query = Query::equal('name', ['John']); $payload = []; - $result = RuntimeQuery::filter([$query], $payload); - $this->assertEquals([], $result); + $result = $this->compileAndFilter([$query], $payload); + $this->assertNull($result); } public function testEmptyAndQuery(): void { $query = Query::and([]); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); // Empty AND should return true (all conditions pass vacuously) $this->assertEquals($payload, $result); } @@ -554,9 +589,9 @@ class RuntimeQueryTest extends TestCase { $query = Query::or([]); $payload = ['name' => 'John']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); // Empty OR should return false (no conditions match) - $this->assertEquals([], $result); + $this->assertNull($result); } // Type-specific edge cases @@ -564,7 +599,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('count', [0]); $payload = ['count' => 0]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -572,7 +607,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('name', ['']); $payload = ['name' => '']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -580,7 +615,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::equal('active', [false]); $payload = ['active' => false]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -588,7 +623,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::greaterThan('score', 8.5); $payload = ['score' => 9.2]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -596,7 +631,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::lessThan('version', '10'); $payload = ['version' => '9']; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -667,7 +702,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::select(['*']); $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } @@ -679,7 +714,7 @@ class RuntimeQueryTest extends TestCase Query::equal('name', ['Jane']), // This would normally fail ]; $payload = ['name' => 'John', 'age' => 30]; - $result = RuntimeQuery::filter($queries, $payload); + $result = $this->compileAndFilter($queries, $payload); // select("*") takes precedence - returns payload $this->assertEquals($payload, $result); } @@ -688,7 +723,7 @@ class RuntimeQueryTest extends TestCase { $query = Query::select(['*']); $payload = []; - $result = RuntimeQuery::filter([$query], $payload); + $result = $this->compileAndFilter([$query], $payload); $this->assertEquals($payload, $result); } }