diff --git a/docs/aggregation.md b/docs/aggregation.md index a198383..6891b57 100644 --- a/docs/aggregation.md +++ b/docs/aggregation.md @@ -17,11 +17,12 @@ foreach ($cursor as $row) { ``` ## Filtering and ordering stages -`$match` filters documents using the same syntax as [query operators](querying.md). `$sort`, `$limit`, and `$skip` order and page the stream just like the [read options](querying.md#sorting): +`$match` filters documents using the same syntax as [query operators](querying.md), including [`$expr`](querying.md#evaluation-operators) for [expression](#expressions)-based conditions. `$sort`, `$limit`, and `$skip` order and page the stream just like the [read options](querying.md#sorting): ```php $collection->aggregate([ ['$match' => ['age' => ['$gte' => 18]]], + ['$match' => ['$expr' => ['$gt' => ['$spent', '$budget']]]], ['$sort' => ['age' => -1]], ['$skip' => 10], ['$limit' => 5], @@ -29,7 +30,7 @@ $collection->aggregate([ ``` ## Reshaping stages -`$project` selects and renames fields, and `$unwind` expands an array field into one document per element: +`$project` selects fields with `1`/`0`, and `$unwind` expands an array field into one document per element: ```php $collection->aggregate([ @@ -40,6 +41,60 @@ $collection->aggregate([ ['$unwind' => '$tags'], ]); ``` + +`$project` can also rename fields and compute new ones from [expressions](#expressions); once a computed field is present the stage rebuilds the document from the keys you list (plus `_id` unless you set `'_id' => 0`): + +```php +$collection->aggregate([ + [ + '$project' => [ + '_id' => 0, + 'fullName' => ['$concat' => ['$first', ' ', '$last']], + 'sku' => '$productId', + 'total' => ['$multiply' => ['$price', '$quantity']], + ], + ], +]); +``` + +`$addFields` (and its alias `$set`) adds or overwrites fields while keeping the rest of the document. Dotted keys write into nested objects: + +```php +$collection->aggregate([ + [ + '$addFields' => [ + 'total' => ['$multiply' => ['$price', '$quantity']], + 'audit.reviewed' => true, + ], + ], +]); +``` + +`$unset` removes one or more fields, and `$replaceRoot` / `$replaceWith` promote an [expression](#expressions) to be the new document: + +```php +$collection->aggregate([ + ['$unset' => ['ssn', 'audit.internalNote']], + ['$replaceRoot' => ['newRoot' => '$profile']], + ['$replaceWith' => ['id' => '$_id', 'name' => '$profile.handle']], +]); +``` + +`$unwind` also accepts the document form with `preserveNullAndEmptyArrays` to keep documents whose array is missing, `null`, or empty, and `includeArrayIndex` to add the position of each element: + +```php +$collection->aggregate([ + [ + '$unwind' => [ + 'path' => '$tags', + 'preserveNullAndEmptyArrays' => true, + 'includeArrayIndex' => 'tagIndex', + ], + ], +]); +``` + +A field that is neither an array nor `null` is treated as a single-element array. ## Grouping `$group` buckets documents by an `_id` expression and computes accumulators per bucket. A field reference is written with a leading `$`: @@ -54,11 +109,81 @@ $collection->aggregate([ 'average' => ['$avg' => '$total'], 'highest' => ['$max' => '$total'], 'lowest' => ['$min' => '$total'], + 'items' => ['$push' => '$sku'], + 'customers' => ['$addToSet' => '$customerId'], ], ], ]); ``` -The supported accumulators are `$sum`, `$avg`, `$min`, `$max`, `$first`, and `$last`. Use `['$sum' => 1]` to count documents in each group. +The supported accumulators are `$sum`, `$avg`, `$min`, `$max`, `$first`, `$last`, `$push`, `$addToSet`, and `$count`. Use `['$sum' => 1]` or `['$count' => []]` to count documents in each group. + +`_id` may also be a document to group by several keys at once, or any [expression](#expressions). Accumulator arguments are expressions too, so you can group over a computed value: + +```php +$collection->aggregate([ + [ + '$group' => [ + '_id' => ['status' => '$status', 'country' => '$address.country'], + 'revenue' => ['$sum' => ['$multiply' => ['$price', '$quantity']]], + ], + ], +]); +``` + +## Expressions + +Wherever a stage expects an expression (`$project`, `$addFields`/`$set`, `$group` keys and accumulator arguments) you can use a field reference (`'$field'`, dot notation allowed), a literal, or an operator object. The supported operators are: + +* Arithmetic: `$add`, `$subtract`, `$multiply`, `$divide`, `$mod`, `$abs`, `$ceil`, `$floor`, `$round` +* String: `$concat`, `$toUpper`, `$toLower`, `$substr`, `$strLenCP` +* Comparison: `$eq`, `$ne`, `$gt`, `$gte`, `$lt`, `$lte` +* Boolean: `$and`, `$or`, `$not` +* Conditional: `$cond`, `$switch`, `$ifNull` +* Type: `$toString`, `$toInt`, `$toLong`, `$toDouble`, `$toBool` +* Date: `$year`, `$month`, `$dayOfMonth`, `$hour`, `$minute`, `$second`, `$dateToString` +* Array: `$size`, `$isArray`, `$arrayElemAt`, `$first`, `$last`, `$in`, `$concatArrays`, `$reverseArray`, `$slice` +* `$literal` to pass a value through untouched + +```php +$collection->aggregate([ + [ + '$project' => [ + 'label' => [ + '$cond' => [ + ['$gte' => ['$score', 60]], + 'pass', + 'fail', + ], + ], + 'month' => ['$dateToString' => ['format' => '%Y-%m', 'date' => '$createdAt']], + ], + ], +]); +``` + +:::note +Date operators expect ISO 8601 date strings, which is how Rango stores dates in JSONB. `$$ROOT` and `$$NOW` are the only system variables. +::: + +## Counting + +`$count` collapses the stream into a single document holding the number of documents that reached it: + +```php +$collection->aggregate([ + ['$match' => ['status' => 'paid']], + ['$count' => 'paidOrders'], +]); +``` + +`$sortByCount` groups by an [expression](#expressions) and returns `{_id, count}` documents ordered by `count` descending. It is shorthand for a `$group` with `['$sum' => 1]` followed by a `$sort`: + +```php +$collection->aggregate([ + ['$unwind' => '$tags'], + ['$sortByCount' => '$tags'], +]); +``` ## Joining collections @@ -79,7 +204,7 @@ $client->selectCollection('app', 'users')->aggregate([ Each `users` document gains an `orders` array holding the matching `orders` documents, or an empty array when there are none. :::note -Only the stages and accumulators listed here are implemented. Complex aggregation expressions are out of scope, as noted under [limitations](how-it-works.md#limitations). +Only the stages, accumulators, and [expression](#expressions) operators listed here are implemented. Array iteration (`$map`, `$filter`, `$reduce`), `$facet`, and window functions are out of scope, as noted under [limitations](how-it-works.md#limitations). ::: ## Learn more diff --git a/docs/how-it-works.md b/docs/how-it-works.md index b298122..8db4215 100644 --- a/docs/how-it-works.md +++ b/docs/how-it-works.md @@ -43,7 +43,7 @@ Rango covers the most common MongoDB use cases, but it does not reimplement the * **Geospatial queries** such as `$near` and `$geoWithin` * **Capped collections** * **Text search** with MongoDB-specific syntax and text indexes -* **Complex aggregation expressions**, beyond the basic accumulators in [aggregation](aggregation.md) +* **Advanced aggregation expressions**: arithmetic, string, comparison, boolean, conditional, date, and array operators are supported (see [aggregation](aggregation.md)), but array iteration (`$map`, `$reduce`, `$filter`), `$facet`, and window functions are not * **Special index types**: only ascending and descending [indexes](indexes.md) are supported, so geospatial (`2dsphere`), text, sparse, and TTL indexes are not, and the matching `IndexInfo` checks always report `false` [Upserts](update-operators.md) also need `_id` to be present in the filter, because Rango builds the primary key of the inserted document from it. An upsert without `_id` in the filter raises an exception. diff --git a/docs/querying.md b/docs/querying.md index 330565b..0997d08 100644 --- a/docs/querying.md +++ b/docs/querying.md @@ -73,6 +73,12 @@ $collection->find(['age' => ['$type' => 'number']]); $collection->find(['email' => ['$regex' => '@example\\.com$']]); $collection->find(['age' => ['$mod' => [2, 0]]]); // even ages ``` + +`$expr` matches documents against an [aggregation expression](aggregation.md#expressions), which is the way to compare two fields of the same document: + +```php +$collection->find(['$expr' => ['$gt' => ['$spent', '$budget']]]); +``` ## Array operators Array operators inspect array fields. `$all` requires every listed value, `$size` matches by length, and `$elemMatch` matches array elements against a sub-filter: diff --git a/src/Query/AggregationBuilder.php b/src/Query/AggregationBuilder.php index c76dba3..83b86b5 100644 --- a/src/Query/AggregationBuilder.php +++ b/src/Query/AggregationBuilder.php @@ -7,15 +7,16 @@ use Patchlevel\Rango\Sql\Identifier; use PDO; +use function array_key_exists; use function explode; use function implode; -use function is_numeric; +use function is_array; +use function is_scalar; use function is_string; -use function json_encode; use function ltrim; use function sprintf; +use function str_contains; use function str_replace; -use function str_starts_with; final readonly class AggregationBuilder { @@ -23,6 +24,7 @@ public function __construct( private PDO $pdo, private FilterBuilder $filterBuilder, private ProjectionBuilder $projectionBuilder, + private ExpressionBuilder $expressionBuilder, ) { } @@ -61,30 +63,54 @@ public function createAggregate(string $database, string $collection, array $pip } elseif ($operator === '$skip') { $currentQuery = sprintf('SELECT * FROM (%s) AS t OFFSET %d', $currentQuery, (int)$value); } elseif ($operator === '$project') { - $column = $this->projectionBuilder->buildProjectionColumn($value, 'data'); - $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $column, $currentQuery); + $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $this->projectColumn($value), $currentQuery); + } elseif ($operator === '$addFields' || $operator === '$set') { + $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $this->addFieldsExpression($value), $currentQuery); + } elseif ($operator === '$unset') { + $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $this->unsetExpression($value), $currentQuery); + } elseif ($operator === '$replaceRoot') { + $newRoot = is_array($value) && array_key_exists('newRoot', $value) ? $value['newRoot'] : $value; + $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $this->expressionBuilder->compile($newRoot), $currentQuery); + } elseif ($operator === '$replaceWith') { + $currentQuery = sprintf('SELECT %s AS data FROM (%s) AS t', $this->expressionBuilder->compile($value), $currentQuery); } elseif ($operator === '$unwind') { - $field = ltrim($value, '$'); + $spec = is_array($value) ? $value : ['path' => $value]; + $field = ltrim($spec['path'], '$'); + $preserve = (bool)($spec['preserveNullAndEmptyArrays'] ?? false); + $includeArrayIndex = $spec['includeArrayIndex'] ?? null; + + $source = $this->fieldReference($field); + $arrayExpr = sprintf( + 'CASE WHEN jsonb_typeof(%1$s) = \'array\' THEN %1$s' + . ' WHEN %1$s IS NULL OR jsonb_typeof(%1$s) = \'null\' THEN \'[]\'::jsonb' + . ' ELSE jsonb_build_array(%1$s) END', + $source, + ); + + $dataExpr = sprintf( + 'CASE WHEN u.elem IS NULL THEN data ELSE jsonb_set(data, %s, u.elem, true) END', + $this->pathLiteral($field), + ); + + if (is_string($includeArrayIndex) && $includeArrayIndex !== '') { + $dataExpr = sprintf( + 'jsonb_set(%s, %s, CASE WHEN u.ord IS NULL THEN \'null\'::jsonb ELSE to_jsonb(u.ord - 1) END, true)', + $dataExpr, + $this->pathLiteral(ltrim($includeArrayIndex, '$')), + ); + } + $currentQuery = sprintf( - 'SELECT jsonb_set(data, %1$s, x) AS data FROM (%2$s) AS t, jsonb_array_elements(CASE WHEN jsonb_typeof(data->%3$s) = \'array\' THEN data->%3$s ELSE \'[]\'::jsonb END) x', - $this->pdo->quote('{' . $field . '}'), + 'SELECT %s AS data FROM (%s) AS t %s LATERAL jsonb_array_elements(%s) WITH ORDINALITY AS u(elem, ord) ON true', + $dataExpr, $currentQuery, - $this->pdo->quote($field), + $preserve ? 'LEFT JOIN' : 'INNER JOIN', + $arrayExpr, ); } elseif ($operator === '$group') { - $id = $value['_id'] ?? null; - $fields = []; - $groupBy = '1'; - if ($id !== null) { - if (is_string($id) && str_starts_with($id, '$')) { - $groupBy = sprintf('data->%s', $this->pdo->quote(ltrim($id, '$'))); - } else { - $groupBy = $this->pdo->quote(json_encode($id)); - } - } + $groupBy = $this->groupKey($value['_id'] ?? null); - $fields[] = $this->pdo->quote('_id'); - $fields[] = $groupBy; + $fields = [$this->pdo->quote('_id'), $groupBy]; foreach ($value as $alias => $expr) { if ($alias === '_id') { @@ -92,29 +118,13 @@ public function createAggregate(string $database, string $collection, array $pip } foreach ($expr as $accOp => $accVal) { - if ($accOp === '$sum') { - $fields[] = $this->pdo->quote($alias); - if (is_numeric($accVal)) { - $fields[] = sprintf('SUM(%s)::text::jsonb', (float)$accVal); - } else { - $fields[] = sprintf('SUM((data->>%s)::numeric)::text::jsonb', $this->pdo->quote(ltrim($accVal, '$'))); - } - } elseif ($accOp === '$avg') { - $fields[] = $this->pdo->quote($alias); - $fields[] = sprintf('AVG((data->>%s)::numeric)::text::jsonb', $this->pdo->quote(ltrim($accVal, '$'))); - } elseif ($accOp === '$min') { - $fields[] = $this->pdo->quote($alias); - $fields[] = sprintf('MIN((data->>%s)::numeric)::text::jsonb', $this->pdo->quote(ltrim($accVal, '$'))); - } elseif ($accOp === '$max') { - $fields[] = $this->pdo->quote($alias); - $fields[] = sprintf('MAX((data->>%s)::numeric)::text::jsonb', $this->pdo->quote(ltrim($accVal, '$'))); - } elseif ($accOp === '$first') { - $fields[] = $this->pdo->quote($alias); - $fields[] = sprintf('(ARRAY_AGG(data->%s))[1]', $this->pdo->quote(ltrim($accVal, '$'))); - } elseif ($accOp === '$last') { - $fields[] = $this->pdo->quote($alias); - $fields[] = sprintf('(ARRAY_AGG(data->%s))[ARRAY_LENGTH(ARRAY_AGG(data->%s), 1)]', $this->pdo->quote(ltrim($accVal, '$')), $this->pdo->quote(ltrim($accVal, '$'))); + $accumulator = $this->accumulator($accOp, $accVal); + if ($accumulator === null) { + continue; } + + $fields[] = $this->pdo->quote($alias); + $fields[] = $accumulator; } } @@ -124,6 +134,22 @@ public function createAggregate(string $database, string $collection, array $pip $currentQuery, $groupBy, ); + } elseif ($operator === '$count') { + $currentQuery = sprintf( + 'SELECT jsonb_build_object(%s, COUNT(*)::text::jsonb) AS data FROM (%s) AS t HAVING COUNT(*) > 0', + $this->pdo->quote($value), + $currentQuery, + ); + } elseif ($operator === '$sortByCount') { + $groupBy = $this->groupKey($value); + $currentQuery = sprintf( + 'SELECT jsonb_build_object(%s, %s, %s, COUNT(*)::text::jsonb) AS data FROM (%s) AS t GROUP BY %s ORDER BY COUNT(*) DESC', + $this->pdo->quote('_id'), + $groupBy, + $this->pdo->quote('count'), + $currentQuery, + $groupBy, + ); } elseif ($operator === '$lookup') { $from = $value['from']; $localField = $value['localField']; @@ -173,4 +199,160 @@ public function createAggregate(string $database, string $collection, array $pip return $currentQuery; } + + private function projectColumn(mixed $projection): string + { + if (!is_array($projection)) { + return 'data'; + } + + $spec = []; + foreach ($projection as $key => $value) { + $spec[(string)$key] = $value; + } + + $hasComputed = false; + foreach ($spec as $value) { + if (is_string($value) || is_array($value)) { + $hasComputed = true; + + break; + } + } + + if (!$hasComputed) { + return $this->projectionBuilder->buildProjectionColumn($spec, 'data'); + } + + $fields = []; + $includeId = true; + + foreach ($spec as $name => $value) { + if ($name === '_id' && ($value === 0 || $value === false)) { + $includeId = false; + + continue; + } + + if ($value === 0 || $value === false) { + continue; + } + + $fields[$name] = $value === 1 || $value === true + ? $this->fieldReference($name) + : $this->expressionBuilder->compile($value); + } + + if ($includeId && !isset($fields['_id'])) { + $fields = ['_id' => "data->'_id'"] + $fields; + } + + $parts = []; + foreach ($fields as $name => $expression) { + $parts[] = $this->pdo->quote($name); + $parts[] = $expression; + } + + return sprintf('jsonb_build_object(%s)', implode(', ', $parts)); + } + + private function addFieldsExpression(mixed $fields): string + { + if (!is_array($fields)) { + return 'data'; + } + + $topLevel = []; + $nested = []; + foreach ($fields as $key => $value) { + $name = (string)$key; + if (str_contains($name, '.')) { + $nested[$name] = $value; + } else { + $topLevel[$name] = $value; + } + } + + $expression = 'data'; + + if ($topLevel !== []) { + $parts = []; + foreach ($topLevel as $name => $value) { + $parts[] = $this->pdo->quote($name); + $parts[] = $this->expressionBuilder->compile($value); + } + + $expression = sprintf('%s || jsonb_build_object(%s)', $expression, implode(', ', $parts)); + } + + foreach ($nested as $name => $value) { + $expression = sprintf( + 'jsonb_set(%s, %s, %s, true)', + $expression, + $this->pathLiteral($name), + $this->expressionBuilder->compile($value), + ); + } + + return $expression; + } + + private function unsetExpression(mixed $fields): string + { + $list = is_array($fields) ? $fields : [$fields]; + + $expression = 'data'; + foreach ($list as $field) { + if (!is_scalar($field)) { + continue; + } + + $name = (string)$field; + if (str_contains($name, '.')) { + $expression = sprintf('(%s) #- %s::text[]', $expression, $this->pathLiteral($name)); + } else { + $expression = sprintf('(%s) - %s', $expression, $this->pdo->quote($name)); + } + } + + return $expression; + } + + private function groupKey(mixed $id): string + { + return $this->expressionBuilder->compile($id); + } + + private function accumulator(mixed $operator, mixed $value): string|null + { + return match ($operator) { + '$sum' => sprintf('COALESCE(SUM(%s), 0)::text::jsonb', $this->expressionBuilder->compileNumeric($value)), + '$avg' => sprintf('AVG(%s)::text::jsonb', $this->expressionBuilder->compileNumeric($value)), + '$min' => sprintf('MIN(%s)::text::jsonb', $this->expressionBuilder->compileNumeric($value)), + '$max' => sprintf('MAX(%s)::text::jsonb', $this->expressionBuilder->compileNumeric($value)), + '$first' => sprintf('(ARRAY_AGG(%s))[1]', $this->expressionBuilder->compile($value)), + '$last' => sprintf( + '(ARRAY_AGG(%1$s))[ARRAY_LENGTH(ARRAY_AGG(%1$s), 1)]', + $this->expressionBuilder->compile($value), + ), + '$push' => sprintf('jsonb_agg(%s)', $this->expressionBuilder->compile($value)), + '$addToSet' => sprintf('jsonb_agg(DISTINCT %s)', $this->expressionBuilder->compile($value)), + '$count' => 'COUNT(*)::text::jsonb', + default => null, + }; + } + + private function fieldReference(string $path): string + { + if (!str_contains($path, '.')) { + return sprintf('data->%s', $this->pdo->quote($path)); + } + + return sprintf('data#>%s', $this->pathLiteral($path)); + } + + private function pathLiteral(string $path): string + { + return $this->pdo->quote('{' . str_replace('.', ',', $path) . '}'); + } } diff --git a/src/Query/ExpressionBuilder.php b/src/Query/ExpressionBuilder.php new file mode 100644 index 0000000..af28c93 --- /dev/null +++ b/src/Query/ExpressionBuilder.php @@ -0,0 +1,467 @@ +systemVariable($expression); + } + + if (str_starts_with($expression, '$')) { + return $this->fieldReference(substr($expression, 1)); + } + + return $this->literal($expression); + } + + if (!is_array($expression)) { + return $this->literal($expression); + } + + if ($expression === [] || array_is_list($expression)) { + $items = array_map(fn (mixed $item): string => $this->compile($item), $expression); + + return sprintf('jsonb_build_array(%s)', implode(', ', $items)); + } + + $firstKey = array_key_first($expression); + if (is_string($firstKey) && str_starts_with($firstKey, '$')) { + if (count($expression) !== 1) { + throw new InvalidArgumentException( + sprintf('Expression object for "%s" must contain exactly one operator', $firstKey), + ); + } + + return $this->operator($firstKey, $expression[$firstKey]); + } + + $parts = []; + foreach ($expression as $key => $value) { + $parts[] = $this->pdo->quote((string)$key); + $parts[] = $this->compile($value); + } + + return sprintf('jsonb_build_object(%s)', implode(', ', $parts)); + } + + /** Compile an expression into an SQL boolean using MongoDB truthiness rules. */ + public function compileBoolean(mixed $expression): string + { + $sql = $this->compile($expression); + + return sprintf( + '(%1$s IS NOT NULL AND %1$s <> \'false\'::jsonb AND %1$s <> \'null\'::jsonb AND %1$s <> \'0\'::jsonb)', + $sql, + ); + } + + /** Compile an expression into a numeric SQL value. */ + public function compileNumeric(mixed $expression): string + { + return sprintf('(%s #>> \'{}\')::numeric', $this->compile($expression)); + } + + private function operator(string $operator, mixed $argument): string + { + return match ($operator) { + '$literal' => $this->literal($argument), + + '$concat' => $this->concat($this->operands($argument)), + '$toUpper' => sprintf('to_jsonb(upper(coalesce(%s, \'\')))', $this->text($argument)), + '$toLower' => sprintf('to_jsonb(lower(coalesce(%s, \'\')))', $this->text($argument)), + '$substr', '$substrCP', '$substrBytes' => $this->substr($this->operands($argument)), + '$strLenCP', '$strLenBytes' => sprintf('to_jsonb(length(coalesce(%s, \'\')))', $this->text($argument)), + '$toString' => sprintf('to_jsonb(%s)', $this->text($argument)), + '$toInt', '$toLong' => sprintf('to_jsonb((%s)::bigint)', $this->text($argument)), + '$toDouble', '$toDecimal' => sprintf('to_jsonb(%s)', $this->compileNumeric($argument)), + '$toBool' => sprintf('to_jsonb(%s)', $this->compileBoolean($argument)), + + '$add' => $this->arithmetic('+', $this->operands($argument)), + '$subtract' => $this->arithmetic('-', $this->operands($argument)), + '$multiply' => $this->arithmetic('*', $this->operands($argument)), + '$divide' => $this->arithmetic('/', $this->operands($argument)), + '$mod' => $this->modulo($this->operands($argument)), + '$abs' => sprintf('to_jsonb(trim_scale(abs(%s)))', $this->compileNumeric($argument)), + '$ceil' => sprintf('to_jsonb(ceil(%s))', $this->compileNumeric($argument)), + '$floor' => sprintf('to_jsonb(floor(%s))', $this->compileNumeric($argument)), + '$round' => $this->round($this->operands($argument)), + + '$eq' => $this->comparison('IS NOT DISTINCT FROM', $this->operands($argument)), + '$ne' => $this->comparison('IS DISTINCT FROM', $this->operands($argument)), + '$gt' => $this->comparison('>', $this->operands($argument)), + '$gte' => $this->comparison('>=', $this->operands($argument)), + '$lt' => $this->comparison('<', $this->operands($argument)), + '$lte' => $this->comparison('<=', $this->operands($argument)), + + '$and' => $this->logical('AND', $this->operands($argument)), + '$or' => $this->logical('OR', $this->operands($argument)), + '$not' => sprintf('to_jsonb(NOT %s)', $this->compileBoolean($this->firstOperand($argument))), + + '$ifNull' => $this->ifNull($this->operands($argument)), + '$cond' => $this->cond($argument), + '$switch' => $this->switchExpression($argument), + + '$year' => $this->datePart('year', $argument), + '$month' => $this->datePart('month', $argument), + '$dayOfMonth' => $this->datePart('day', $argument), + '$hour' => $this->datePart('hour', $argument), + '$minute' => $this->datePart('minute', $argument), + '$second' => $this->datePart('second', $argument), + '$dateToString' => $this->dateToString($argument), + + '$size' => sprintf( + 'to_jsonb(CASE WHEN jsonb_typeof(%1$s) = \'array\' THEN jsonb_array_length(%1$s) END)', + $this->compile($this->firstOperand($argument)), + ), + '$isArray' => sprintf( + 'to_jsonb(COALESCE(jsonb_typeof(%s) = \'array\', false))', + $this->compile($this->firstOperand($argument)), + ), + '$arrayElemAt' => $this->arrayElemAt($this->operands($argument)), + '$first' => sprintf('(%s) -> 0', $this->compile($this->firstOperand($argument))), + '$last' => sprintf('(%s) -> -1', $this->compile($this->firstOperand($argument))), + '$in' => $this->inArray($this->operands($argument)), + '$concatArrays' => $this->concatArrays($this->operands($argument)), + '$reverseArray' => $this->reverseArray($this->firstOperand($argument)), + '$slice' => $this->slice($this->operands($argument)), + + default => throw new InvalidArgumentException( + sprintf('Unsupported aggregation expression operator "%s"', $operator), + ), + }; + } + + /** @return list */ + private function operands(mixed $argument): array + { + return is_array($argument) && array_is_list($argument) ? $argument : [$argument]; + } + + private function firstOperand(mixed $argument): mixed + { + return $this->operands($argument)[0] ?? null; + } + + private function literal(mixed $value): string + { + return sprintf('%s::jsonb', $this->pdo->quote((string)json_encode($value))); + } + + private function systemVariable(string $expression): string + { + return match ($expression) { + '$$ROOT', '$$CURRENT' => 'data', + '$$NOW' => 'to_jsonb(now())', + default => throw new InvalidArgumentException( + sprintf('Unsupported system variable "%s"', $expression), + ), + }; + } + + private function fieldReference(string $path): string + { + if ($path === '') { + return 'data'; + } + + if (!str_contains($path, '.')) { + return sprintf('data->%s', $this->pdo->quote($path)); + } + + return sprintf('data#>%s', $this->pdo->quote('{' . str_replace('.', ',', $path) . '}')); + } + + private function text(mixed $expression): string + { + return sprintf('(%s #>> \'{}\')', $this->compile($expression)); + } + + /** @param list $operands */ + private function concat(array $operands): string + { + $parts = array_map(fn (mixed $operand): string => $this->text($operand), $operands); + + return sprintf('to_jsonb(%s)', implode(' || ', $parts)); + } + + /** @param list $operands */ + private function substr(array $operands): string + { + $string = $this->text($operands[0] ?? null); + $start = sprintf('(%s)::int', $this->compileNumeric($operands[1] ?? 0)); + $length = sprintf('(%s)::int', $this->compileNumeric($operands[2] ?? -1)); + + return sprintf( + 'to_jsonb(CASE WHEN %3$s < 0 THEN substr(%1$s, %2$s + 1) ELSE substr(%1$s, %2$s + 1, %3$s) END)', + $string, + $start, + $length, + ); + } + + /** @param list $operands */ + private function arithmetic(string $operator, array $operands): string + { + $parts = array_map(fn (mixed $operand): string => $this->compileNumeric($operand), $operands); + + return sprintf('to_jsonb(trim_scale(%s))', implode(sprintf(' %s ', $operator), $parts)); + } + + /** @param list $operands */ + private function modulo(array $operands): string + { + return sprintf( + 'to_jsonb(trim_scale(mod(%s, %s)))', + $this->compileNumeric($operands[0] ?? 0), + $this->compileNumeric($operands[1] ?? 1), + ); + } + + /** @param list $operands */ + private function round(array $operands): string + { + $place = isset($operands[1]) ? sprintf('(%s)::int', $this->compileNumeric($operands[1])) : '0'; + + return sprintf('to_jsonb(trim_scale(round(%s, %s)))', $this->compileNumeric($operands[0] ?? 0), $place); + } + + /** @param list $operands */ + private function comparison(string $operator, array $operands): string + { + return sprintf( + 'to_jsonb((%s) %s (%s))', + $this->compile($operands[0] ?? null), + $operator, + $this->compile($operands[1] ?? null), + ); + } + + /** @param list $operands */ + private function logical(string $operator, array $operands): string + { + if ($operands === []) { + return $operator === 'AND' ? "'true'::jsonb" : "'false'::jsonb"; + } + + $parts = array_map(fn (mixed $operand): string => $this->compileBoolean($operand), $operands); + + return sprintf('to_jsonb(%s)', implode(sprintf(' %s ', $operator), $parts)); + } + + /** @param list $operands */ + private function ifNull(array $operands): string + { + $parts = array_map(fn (mixed $operand): string => $this->compile($operand), $operands); + + return sprintf('COALESCE(%s)', implode(', ', $parts)); + } + + private function cond(mixed $argument): string + { + if (is_array($argument) && array_is_list($argument)) { + [$if, $then, $else] = [$argument[0] ?? null, $argument[1] ?? null, $argument[2] ?? null]; + } elseif (is_array($argument)) { + [$if, $then, $else] = [$argument['if'] ?? null, $argument['then'] ?? null, $argument['else'] ?? null]; + } else { + throw new InvalidArgumentException('$cond expects an array or an object with if/then/else'); + } + + return sprintf( + 'CASE WHEN %s THEN %s ELSE %s END', + $this->compileBoolean($if), + $this->compile($then), + $this->compile($else), + ); + } + + private function switchExpression(mixed $argument): string + { + if (!is_array($argument)) { + throw new InvalidArgumentException('$switch expects an object'); + } + + $branches = $argument['branches'] ?? []; + $default = array_key_exists('default', $argument) + ? $this->compile($argument['default']) + : "'null'::jsonb"; + + if (!is_array($branches) || $branches === []) { + return $default; + } + + $cases = []; + foreach ($branches as $branch) { + if (!is_array($branch)) { + throw new InvalidArgumentException('$switch branch must be an object'); + } + + $cases[] = sprintf( + 'WHEN %s THEN %s', + $this->compileBoolean($branch['case'] ?? null), + $this->compile($branch['then'] ?? null), + ); + } + + return sprintf('CASE %s ELSE %s END', implode(' ', $cases), $default); + } + + private function datePart(string $part, mixed $argument): string + { + $date = $argument; + if (is_array($argument) && !array_is_list($argument) && array_key_exists('date', $argument)) { + $date = $argument['date']; + } + + return sprintf('to_jsonb(extract(%s from (%s)::timestamptz)::int)', $part, $this->text($date)); + } + + private function dateToString(mixed $argument): string + { + if (!is_array($argument)) { + throw new InvalidArgumentException('$dateToString expects an object'); + } + + $format = $argument['format'] ?? '%Y-%m-%dT%H:%M:%S'; + $pgFormat = strtr(is_string($format) ? $format : '%Y-%m-%dT%H:%M:%S', [ + '%Y' => 'YYYY', + '%m' => 'MM', + '%d' => 'DD', + '%H' => 'HH24', + '%M' => 'MI', + '%S' => 'SS', + '%L' => 'MS', + '%j' => 'DDD', + '%%' => '%', + ]); + + return sprintf( + 'to_jsonb(to_char((%s)::timestamptz, %s))', + $this->text($argument['date'] ?? null), + $this->pdo->quote($pgFormat), + ); + } + + /** @param list $operands */ + private function arrayElemAt(array $operands): string + { + return sprintf( + '(%s) -> %s', + $this->compile($operands[0] ?? null), + $this->intExpression($operands[1] ?? 0), + ); + } + + /** @param list $operands */ + private function inArray(array $operands): string + { + return sprintf( + 'to_jsonb(EXISTS (SELECT 1 FROM jsonb_array_elements(%s) AS __in(__v) WHERE __in.__v = %s))', + $this->arrayCase($operands[1] ?? null), + $this->compile($operands[0] ?? null), + ); + } + + /** @param list $operands */ + private function concatArrays(array $operands): string + { + if ($operands === []) { + return "'[]'::jsonb"; + } + + $parts = array_map(fn (mixed $operand): string => sprintf('(%s)', $this->compile($operand)), $operands); + + return implode(' || ', $parts); + } + + private function reverseArray(mixed $argument): string + { + $sql = $this->compile($argument); + + return sprintf( + 'CASE WHEN jsonb_typeof(%1$s) = \'array\' THEN (' + . 'SELECT COALESCE(jsonb_agg(__r.__v ORDER BY __r.__ord DESC), \'[]\'::jsonb)' + . ' FROM jsonb_array_elements(%1$s) WITH ORDINALITY AS __r(__v, __ord)) END', + $sql, + ); + } + + /** @param list $operands */ + private function slice(array $operands): string + { + $array = $this->arrayCase($operands[0] ?? null); + + if (array_key_exists(2, $operands)) { + $position = $this->intExpression($operands[1] ?? 0); + $count = $this->intExpression($operands[2] ?? 0); + $condition = sprintf( + 'CASE WHEN %1$s >= 0 THEN q.__ord > %1$s AND q.__ord <= %1$s + %2$s' + . ' ELSE q.__ord > q.__total + %1$s AND q.__ord <= q.__total + %1$s + %2$s END', + $position, + $count, + ); + } else { + $count = $this->intExpression($operands[1] ?? 0); + $condition = sprintf( + 'CASE WHEN %1$s >= 0 THEN q.__ord <= %1$s ELSE q.__ord > q.__total + %1$s END', + $count, + ); + } + + return sprintf( + '(SELECT COALESCE(jsonb_agg(q.__v ORDER BY q.__ord), \'[]\'::jsonb) FROM (' + . 'SELECT __s.__v AS __v, __s.__ord AS __ord, count(*) OVER () AS __total' + . ' FROM jsonb_array_elements(%s) WITH ORDINALITY AS __s(__v, __ord)) q WHERE %s)', + $array, + $condition, + ); + } + + private function arrayCase(mixed $expression): string + { + $sql = $this->compile($expression); + + return sprintf( + 'CASE WHEN jsonb_typeof(%1$s) = \'array\' THEN %1$s ELSE \'[]\'::jsonb END', + $sql, + ); + } + + private function intExpression(mixed $expression): string + { + return sprintf('(%s)::int', $this->compileNumeric($expression)); + } +} diff --git a/src/Query/FilterBuilder.php b/src/Query/FilterBuilder.php index 2fe593a..acc062a 100644 --- a/src/Query/FilterBuilder.php +++ b/src/Query/FilterBuilder.php @@ -21,10 +21,14 @@ final readonly class FilterBuilder { + private ExpressionBuilder $expressionBuilder; + public function __construct( private PDO $pdo, private string|null $idColumn = null, + ExpressionBuilder|null $expressionBuilder = null, ) { + $this->expressionBuilder = $expressionBuilder ?? new ExpressionBuilder($pdo); } /** @param array $filter */ @@ -62,6 +66,11 @@ public function buildFilter(array $filter, string $conjunction = 'AND', string $ continue; } + if ($key === '$expr') { + $parts[] = '(' . $this->expressionBuilder->compileBoolean($value) . ')'; + continue; + } + if (str_starts_with($key, '$')) { throw new RuntimeException(sprintf('Operator "%s" is not supported at the top level', $key)); } diff --git a/src/QueryBuilder.php b/src/QueryBuilder.php index 56a0e06..bf25ef6 100644 --- a/src/QueryBuilder.php +++ b/src/QueryBuilder.php @@ -5,6 +5,7 @@ namespace Patchlevel\Rango; use Patchlevel\Rango\Query\AggregationBuilder; +use Patchlevel\Rango\Query\ExpressionBuilder; use Patchlevel\Rango\Query\FilterBuilder; use Patchlevel\Rango\Query\ProjectionBuilder; use Patchlevel\Rango\Query\UpdateBuilder; @@ -32,10 +33,16 @@ public function __construct( private PDO $pdo, ) { - $this->filterBuilder = new FilterBuilder($pdo, '_id'); + $expressionBuilder = new ExpressionBuilder($pdo); + $this->filterBuilder = new FilterBuilder($pdo, '_id', $expressionBuilder); $this->projectionBuilder = new ProjectionBuilder($pdo); $this->updateBuilder = new UpdateBuilder($pdo); - $this->aggregationBuilder = new AggregationBuilder($pdo, new FilterBuilder($pdo), $this->projectionBuilder); + $this->aggregationBuilder = new AggregationBuilder( + $pdo, + new FilterBuilder($pdo, null, $expressionBuilder), + $this->projectionBuilder, + $expressionBuilder, + ); } /** @param array $document */ diff --git a/tests/IntegrationTest.php b/tests/IntegrationTest.php index 83b8a2a..710c7bc 100644 --- a/tests/IntegrationTest.php +++ b/tests/IntegrationTest.php @@ -20,6 +20,7 @@ use function array_values; use function is_array; use function iterator_to_array; +use function json_decode; use function json_encode; use function sort; @@ -983,6 +984,453 @@ public function testComplexAggregate(): void self::assertEquals(10, $docs[2]['total']); } + public function testUnwindWithIncludeArrayIndex(): void + { + $this->collection->insertOne(['_id' => '1', 'items' => ['a', 'b', 'c']]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$unwind' => ['path' => '$items', 'includeArrayIndex' => 'idx']], + ])); + + self::assertCount(3, $docs); + self::assertSame('a', $docs[0]['items']); + self::assertSame(0, $docs[0]['idx']); + self::assertSame('b', $docs[1]['items']); + self::assertSame(1, $docs[1]['idx']); + self::assertSame('c', $docs[2]['items']); + self::assertSame(2, $docs[2]['idx']); + } + + public function testUnwindPreserveNullAndEmptyArrays(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'items' => ['a', 'b']], + ['_id' => '2', 'items' => []], + ['_id' => '3'], + ['_id' => '4', 'items' => null], + ]); + + $preserved = $this->toPlainArrays($this->collection->aggregate([ + ['$unwind' => ['path' => '$items', 'preserveNullAndEmptyArrays' => true]], + ['$sort' => ['_id' => 1]], + ])); + self::assertSame( + ['1', '1', '2', '3', '4'], + array_map(static fn (array $doc) => $doc['_id'], $preserved), + ); + + $default = $this->toPlainArrays($this->collection->aggregate([ + ['$unwind' => '$items'], + ['$sort' => ['_id' => 1]], + ])); + self::assertSame( + ['1', '1'], + array_map(static fn (array $doc) => $doc['_id'], $default), + ); + } + + public function testUnwindScalarField(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'item' => 'solo'], + ['_id' => '2', 'item' => ['x', 'y']], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$unwind' => '$item'], + ['$sort' => ['_id' => 1]], + ])); + + self::assertCount(3, $docs); + self::assertSame('solo', $docs[0]['item']); + self::assertSame('x', $docs[1]['item']); + self::assertSame('y', $docs[2]['item']); + } + + public function testGroupByMultipleKeys(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'store' => 'a', 'type' => 'x', 'amount' => 10], + ['_id' => '2', 'store' => 'a', 'type' => 'x', 'amount' => 5], + ['_id' => '3', 'store' => 'a', 'type' => 'y', 'amount' => 20], + ['_id' => '4', 'store' => 'b', 'type' => 'x', 'amount' => 7], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$group' => [ + '_id' => ['store' => '$store', 'type' => '$type'], + 'total' => ['$sum' => '$amount'], + ], + ], + ])); + + $rows = []; + foreach ($docs as $doc) { + $id = $doc['_id']; + self::assertIsArray($id); + $rows[] = ['store' => $id['store'], 'type' => $id['type'], 'total' => $doc['total']]; + } + + self::assertEqualsCanonicalizing([ + ['store' => 'a', 'type' => 'x', 'total' => 15], + ['store' => 'a', 'type' => 'y', 'total' => 20], + ['store' => 'b', 'type' => 'x', 'total' => 7], + ], $rows); + } + + public function testGroupPushAndAddToSet(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'cat' => 'A', 'tag' => 'x'], + ['_id' => '2', 'cat' => 'A', 'tag' => 'y'], + ['_id' => '3', 'cat' => 'A', 'tag' => 'x'], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$group' => [ + '_id' => '$cat', + 'all' => ['$push' => '$tag'], + 'unique' => ['$addToSet' => '$tag'], + ], + ], + ])); + + self::assertCount(1, $docs); + + $all = $docs[0]['all']; + self::assertIsArray($all); + sort($all); + self::assertSame(['x', 'x', 'y'], $all); + + $unique = $docs[0]['unique']; + self::assertIsArray($unique); + sort($unique); + self::assertSame(['x', 'y'], $unique); + } + + public function testCountStage(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'status' => 'A'], + ['_id' => '2', 'status' => 'B'], + ['_id' => '3', 'status' => 'A'], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$match' => ['status' => 'A']], + ['$count' => 'total'], + ])); + + self::assertCount(1, $docs); + self::assertSame(2, $docs[0]['total']); + + $empty = $this->toPlainArrays($this->collection->aggregate([ + ['$match' => ['status' => 'Z']], + ['$count' => 'total'], + ])); + self::assertCount(0, $empty); + } + + public function testProjectComputedFields(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'first' => 'Ada', 'last' => 'Lovelace', 'price' => 10, 'qty' => 3], + ['_id' => '2', 'first' => 'Alan', 'last' => 'Turing', 'price' => 5, 'qty' => 4], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$sort' => ['_id' => 1]], + [ + '$project' => [ + '_id' => 0, + 'name' => ['$concat' => ['$first', ' ', '$last']], + 'renamed' => '$first', + 'total' => ['$multiply' => ['$price', '$qty']], + ], + ], + ])); + + self::assertEquals([ + ['name' => 'Ada Lovelace', 'renamed' => 'Ada', 'total' => 30], + ['name' => 'Alan Turing', 'renamed' => 'Alan', 'total' => 20], + ], $docs); + } + + public function testProjectConditional(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'score' => 80], + ['_id' => '2', 'score' => 40], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$sort' => ['_id' => 1]], + [ + '$project' => [ + 'grade' => ['$cond' => [['$gte' => ['$score', 60]], 'pass', 'fail']], + ], + ], + ])); + + self::assertSame('1', $docs[0]['_id']); + self::assertSame('pass', $docs[0]['grade']); + self::assertSame('fail', $docs[1]['grade']); + } + + public function testAddFieldsStage(): void + { + $this->collection->insertOne([ + '_id' => '1', + 'price' => 10, + 'qty' => 3, + 'meta' => ['tag' => 'x'], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$addFields' => [ + 'total' => ['$multiply' => ['$price', '$qty']], + 'meta.checked' => true, + ], + ], + ])); + + self::assertCount(1, $docs); + $doc = $docs[0]; + self::assertSame(10, $doc['price']); + self::assertEquals(30, $doc['total']); + $meta = $doc['meta']; + self::assertIsArray($meta); + self::assertSame('x', $meta['tag']); + self::assertTrue($meta['checked']); + + $viaSet = $this->toPlainArrays($this->collection->aggregate([ + ['$set' => ['doubled' => ['$multiply' => ['$price', 2]]]], + ])); + self::assertEquals(20, $viaSet[0]['doubled']); + } + + public function testGroupSumOfExpression(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'cat' => 'A', 'price' => 10, 'qty' => 2], + ['_id' => '2', 'cat' => 'A', 'price' => 5, 'qty' => 4], + ['_id' => '3', 'cat' => 'B', 'price' => 3, 'qty' => 3], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$group' => [ + '_id' => '$cat', + 'revenue' => ['$sum' => ['$multiply' => ['$price', '$qty']]], + ], + ], + ['$sort' => ['_id' => 1]], + ])); + + self::assertSame('A', $docs[0]['_id']); + self::assertEquals(40, $docs[0]['revenue']); + self::assertSame('B', $docs[1]['_id']); + self::assertEquals(9, $docs[1]['revenue']); + } + + public function testGroupByExpression(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'cat' => 'foo', 'n' => 1], + ['_id' => '2', 'cat' => 'FOO', 'n' => 2], + ['_id' => '3', 'cat' => 'bar', 'n' => 4], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$group' => [ + '_id' => ['$toUpper' => '$cat'], + 'total' => ['$sum' => '$n'], + ], + ], + ['$sort' => ['_id' => 1]], + ])); + + self::assertSame('BAR', $docs[0]['_id']); + self::assertEquals(4, $docs[0]['total']); + self::assertSame('FOO', $docs[1]['_id']); + self::assertEquals(3, $docs[1]['total']); + } + + public function testExprInFind(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'spent' => 120, 'budget' => 100], + ['_id' => '2', 'spent' => 80, 'budget' => 100], + ['_id' => '3', 'spent' => 100, 'budget' => 100], + ]); + + $docs = $this->toPlainArrays($this->collection->find([ + '$expr' => ['$gt' => ['$spent', '$budget']], + ])); + $ids = array_map(static fn (array $doc) => $doc['_id'], $docs); + sort($ids); + + self::assertSame(['1'], $ids); + } + + public function testExprInMatch(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'a' => 5, 'b' => 2], + ['_id' => '2', 'a' => 1, 'b' => 4], + ['_id' => '3', 'a' => 3, 'b' => 3], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$match' => ['$expr' => ['$gte' => ['$a', '$b']]]], + ['$sort' => ['_id' => 1]], + ])); + + self::assertSame(['1', '3'], array_map(static fn (array $doc) => $doc['_id'], $docs)); + } + + public function testUnsetStage(): void + { + $this->collection->insertOne([ + '_id' => '1', + 'name' => 'foo', + 'secret' => 'hide', + 'meta' => ['keep' => 1, 'drop' => 2], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$unset' => ['secret', 'meta.drop']], + ])); + + $doc = $docs[0]; + self::assertSame('foo', $doc['name']); + self::assertArrayNotHasKey('secret', $doc); + $meta = $doc['meta']; + self::assertIsArray($meta); + self::assertSame(1, $meta['keep']); + self::assertArrayNotHasKey('drop', $meta); + + $single = $this->toPlainArrays($this->collection->aggregate([ + ['$unset' => 'name'], + ])); + self::assertArrayNotHasKey('name', $single[0]); + } + + public function testReplaceRootStage(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'profile' => ['handle' => 'ada', 'age' => 30]], + ['_id' => '2', 'profile' => ['handle' => 'alan', 'age' => 40]], + ]); + + $viaRoot = $this->toPlainArrays($this->collection->aggregate([ + ['$sort' => ['_id' => 1]], + ['$replaceRoot' => ['newRoot' => '$profile']], + ])); + self::assertSame('ada', $viaRoot[0]['handle']); + self::assertSame(30, $viaRoot[0]['age']); + self::assertArrayNotHasKey('profile', $viaRoot[0]); + + $viaWith = $this->toPlainArrays($this->collection->aggregate([ + ['$sort' => ['_id' => 1]], + ['$replaceWith' => ['id' => '$_id', 'name' => '$profile.handle']], + ])); + self::assertSame('1', $viaWith[0]['id']); + self::assertSame('ada', $viaWith[0]['name']); + } + + public function testSortByCountStage(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'type' => 'a'], + ['_id' => '2', 'type' => 'b'], + ['_id' => '3', 'type' => 'a'], + ['_id' => '4', 'type' => 'a'], + ['_id' => '5', 'type' => 'b'], + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + ['$sortByCount' => '$type'], + ])); + + self::assertCount(2, $docs); + self::assertSame('a', $docs[0]['_id']); + self::assertEquals(3, $docs[0]['count']); + self::assertSame('b', $docs[1]['_id']); + self::assertEquals(2, $docs[1]['count']); + } + + public function testArrayExpressionOperators(): void + { + $this->collection->insertOne([ + '_id' => '1', + 'nums' => [10, 20, 30, 40], + 'tags' => ['a', 'b', 'c'], + 'more' => ['d', 'e'], + 'scalar' => 5, + ]); + + $docs = $this->toPlainArrays($this->collection->aggregate([ + [ + '$project' => [ + '_id' => 0, + 'size' => ['$size' => '$nums'], + 'firstNum' => ['$first' => '$nums'], + 'lastNum' => ['$last' => '$nums'], + 'third' => ['$arrayElemAt' => ['$nums', 2]], + 'fromEnd' => ['$arrayElemAt' => ['$nums', -1]], + 'hasB' => ['$in' => ['b', '$tags']], + 'hasZ' => ['$in' => ['z', '$tags']], + 'numsIsArray' => ['$isArray' => '$nums'], + 'scalarIsArray' => ['$isArray' => '$scalar'], + 'combined' => ['$concatArrays' => ['$tags', '$more']], + 'reversed' => ['$reverseArray' => '$tags'], + 'firstTwo' => ['$slice' => ['$nums', 2]], + 'lastTwo' => ['$slice' => ['$nums', -2]], + 'middle' => ['$slice' => ['$nums', 1, 2]], + ], + ], + ])); + + $doc = $docs[0]; + self::assertSame(4, $doc['size']); + self::assertSame(10, $doc['firstNum']); + self::assertSame(40, $doc['lastNum']); + self::assertSame(30, $doc['third']); + self::assertSame(40, $doc['fromEnd']); + self::assertTrue($doc['hasB']); + self::assertFalse($doc['hasZ']); + self::assertTrue($doc['numsIsArray']); + self::assertFalse($doc['scalarIsArray']); + self::assertSame(['a', 'b', 'c', 'd', 'e'], $doc['combined']); + self::assertSame(['c', 'b', 'a'], $doc['reversed']); + self::assertSame([10, 20], $doc['firstTwo']); + self::assertSame([30, 40], $doc['lastTwo']); + self::assertSame([20, 30], $doc['middle']); + } + + /** + * @param iterable $result + * + * @return list> + */ + protected function toPlainArrays(iterable $result): array + { + $documents = []; + + foreach ($result as $document) { + $decoded = json_decode((string)json_encode($document), true); + self::assertIsArray($decoded); + $documents[] = $decoded; + } + + return $documents; + } + public function testCombinedUpdate(): void { $this->collection->insertOne([ diff --git a/tests/PostgresIntegrationTest.php b/tests/PostgresIntegrationTest.php index e695c9b..a54fb31 100644 --- a/tests/PostgresIntegrationTest.php +++ b/tests/PostgresIntegrationTest.php @@ -47,4 +47,47 @@ protected function getDatabase(): Database { return $this->getClient()->getDatabase('test'); } + + /** + * Date expression operators work on ISO 8601 date strings, which is how + * Rango stores dates in JSONB. Real MongoDB requires a BSON date here, so + * this behaviour is verified against Postgres only. + */ + public function testDateExpressionOperators(): void + { + $this->collection->insertMany([ + ['_id' => '1', 'ts' => '2024-03-15T12:00:00Z', 'amount' => 10], + ['_id' => '2', 'ts' => '2024-03-20T12:00:00Z', 'amount' => 5], + ['_id' => '3', 'ts' => '2024-04-01T12:00:00Z', 'amount' => 7], + ]); + + $byMonth = $this->toPlainArrays($this->collection->aggregate([ + [ + '$group' => [ + '_id' => ['$dateToString' => ['format' => '%Y-%m', 'date' => '$ts']], + 'revenue' => ['$sum' => '$amount'], + ], + ], + ['$sort' => ['_id' => 1]], + ])); + + self::assertSame('2024-03', $byMonth[0]['_id']); + self::assertEquals(15, $byMonth[0]['revenue']); + self::assertSame('2024-04', $byMonth[1]['_id']); + self::assertEquals(7, $byMonth[1]['revenue']); + + $parts = $this->toPlainArrays($this->collection->aggregate([ + ['$match' => ['_id' => '1']], + [ + '$project' => [ + '_id' => 0, + 'year' => ['$year' => '$ts'], + 'month' => ['$month' => '$ts'], + 'day' => ['$dayOfMonth' => '$ts'], + ], + ], + ])); + + self::assertEquals(['year' => 2024, 'month' => 3, 'day' => 15], $parts[0]); + } }