From 5219359adac2aa004c01b97020ef0f780ce73525 Mon Sep 17 00:00:00 2001 From: Jan Kaniuka Date: Sat, 29 Aug 2026 19:08:34 +0200 Subject: [PATCH 1/2] feat: CQL2 support for Mongo provider --- pygeoapi/provider/mongo.py | 133 +++++++++++++++++++++++++++++++------ 1 file changed, 112 insertions(+), 21 deletions(-) diff --git a/pygeoapi/provider/mongo.py b/pygeoapi/provider/mongo.py index 049174f3f9..698260ae28 100644 --- a/pygeoapi/provider/mongo.py +++ b/pygeoapi/provider/mongo.py @@ -66,7 +66,7 @@ def __init__(self, provider_def): dbclient = MongoClient(self.data) self.featuredb = dbclient.get_default_database() self.collection = provider_def['collection'] - self.featuredb[self.collection].create_index([("geometry", GEOSPHERE)]) + self.featuredb[self.collection].create_index([("feature.geometry", GEOSPHERE)]) self.get_fields() def get_fields(self): @@ -98,45 +98,142 @@ def get_fields(self): def _get_feature_list(self, filterObj, sortList=[], skip=0, maxitems=1, skip_geometry=False): featurecursor = self.featuredb[self.collection].find(filterObj) - if sortList: featurecursor = featurecursor.sort(sortList) featurecursor.skip(skip) - featurecursor.limit(maxitems) + if maxitems > -1: + featurecursor.limit(maxitems) featurelist = list(featurecursor) + + features = [] + for item in featurelist: - item['id'] = str(item.pop('_id')) + feature_id = str(item.pop('_id')) + geometry = item['feature']['geometry'] + props = item['feature']['properties'] + if skip_geometry: - item['geometry'] = None + geometry = None - return featurelist + feature = { + 'type': 'Feature', + 'id': feature_id, + 'geometry': geometry, + 'properties': props + } + + features.append(feature) + + return features @crs_transform def query(self, offset=0, limit=10, resulttype='results', bbox=[], datetime_=None, properties=[], sortby=[], - select_properties=[], skip_geometry=False, q=None, **kwargs): + select_properties=[], skip_geometry=False, q=None, filterq=None, **kwargs): """ query the provider :returns: dict of 0..n GeoJSON features """ - and_filter = [] + def cql2_to_mongo(node): + if node is None: + return + + # GeoJson operator + if node.__class__.__name__ == "GeometryWithin": + field = node.lhs.name + geom = node.rhs.geometry + + query_body = { + field: { + '$geoWithin': { + '$geometry': geom + } + } + } + return query_body + + if node.__class__.__name__ == "GeometryIntersects": + field = node.lhs.name + geom = node.rhs.geometry + + return { + field: { + '$geoIntersects': { + '$geometry': geom + } + } + } + + if node.__class__.__name__ == "DistanceWithin": + field = node.lhs.name + geom = node.rhs.geometry + distance = node.distance + units = node.units # mongo's default units are meters + # but with CQL we can pass different units + # and here we can recalculate them + return { + field: { + '$near': { + '$geometry': geom, + '$maxDistance': distance, + '$minDistance': 0 + } + } + } + + # Logical operators + if node.__class__.__name__ == "And": + return { + "$and": [ + cql2_to_mongo(node.lhs), + cql2_to_mongo(node.rhs) + ] + } + + if node.__class__.__name__ == "Or": + return { + "$or": [ + cql2_to_mongo(node.lhs), + cql2_to_mongo(node.rhs) + ] + } + + # Comparison operators + if node.__class__.__name__ == "Equal": + field = node.lhs.name + value = node.rhs + return {f"{field}": value} + + if node.__class__.__name__ == "GreaterEqual": + field = node.lhs.name + value = node.rhs + return {f"{field}": {"$gte": value}} + + if node.__class__.__name__ == "LessEqual": + field = node.lhs.name + value = node.rhs + return {f"{field}": {"$lte": value}} + + return + + and_filter = [] + cql_filters_parsed = cql2_to_mongo(filterq) + if cql_filters_parsed is not None: + and_filter.append(cql_filters_parsed) + limit = -1 # if there is CQL query return all elements + if len(bbox) == 4: x, y, w, h = map(float, bbox) and_filter.append( {'geometry': {'$geoWithin': {'$box': [[x, y], [w, h]]}}}) - # This parameter is not working yet! - # gte is not sufficient to check date range - if datetime_ is not None: - assert isinstance(datetime_, datetime) - and_filter.append({'properties.datetime': {'$gte': datetime_}}) - + for prop in properties: and_filter.append({"properties."+prop[0]: {'$eq': prop[1]}}) - + filterobj = {'$and': and_filter} if and_filter else {} sort_list = [("properties." + sort['property'], @@ -148,12 +245,6 @@ def query(self, offset=0, limit=10, resulttype='results', 'features': [] } - if self.count or resulttype == 'hits': - matched = self.featuredb[self.collection].count_documents( - filterobj) - LOGGER.debug(f'Found {matched} result(s)') - feature_collection['numberMatched'] = matched - if resulttype == 'hits': return feature_collection From caff9768a3fc2872c5035d68419597a74e5cb54c Mon Sep 17 00:00:00 2001 From: Jan Kaniuka Date: Tue, 15 Sep 2026 15:53:44 +0200 Subject: [PATCH 2/2] fix: pre-commit --- pygeoapi/api/collection.py | 2 +- pygeoapi/provider/mongo.py | 172 +++++++++++++++++++------------------ 2 files changed, 90 insertions(+), 84 deletions(-) diff --git a/pygeoapi/api/collection.py b/pygeoapi/api/collection.py index 2cd8b3ea62..584490a629 100644 --- a/pygeoapi/api/collection.py +++ b/pygeoapi/api/collection.py @@ -437,7 +437,7 @@ def gen_collection(api, request, dataset: str, 'type': FORMAT_TYPES[F_HTML], 'rel': 'data', 'title': title2, - 'href': f'{api.get_collections_url()}/{dataset}/{qt}?f={F_HTML}' # noqa + 'href': f'{api.get_collections_url()}/{dataset}/{qt}?f={F_HTML}' # noqa }]) for key, value in get_dataset_formatters(config).items(): diff --git a/pygeoapi/provider/mongo.py b/pygeoapi/provider/mongo.py index 698260ae28..ed064de5dc 100644 --- a/pygeoapi/provider/mongo.py +++ b/pygeoapi/provider/mongo.py @@ -28,7 +28,6 @@ # # ================================================================= -from datetime import datetime import logging from pymongo import MongoClient @@ -43,8 +42,7 @@ class MongoProvider(BaseProvider): - """Generic provider for Mongodb. - """ + """Generic provider for Mongodb.""" def __init__(self, provider_def): """ @@ -57,16 +55,18 @@ def __init__(self, provider_def): """ # this is dummy value never used in case of Mongo. # Mongo id field is _id - provider_def.setdefault('id_field', '_id') + provider_def.setdefault("id_field", "_id") super().__init__(provider_def) - LOGGER.info(f'Mongo source config: {self.data}') + LOGGER.info(f"Mongo source config: {self.data}") dbclient = MongoClient(self.data) self.featuredb = dbclient.get_default_database() - self.collection = provider_def['collection'] - self.featuredb[self.collection].create_index([("feature.geometry", GEOSPHERE)]) + self.collection = provider_def["collection"] + self.featuredb[self.collection].create_index( + [("feature.geometry", GEOSPHERE)] + ) self.get_fields() def get_fields(self): @@ -81,7 +81,7 @@ def get_fields(self): {"$project": {"properties": 1}}, {"$unwind": "$properties"}, {"$group": {"_id": "$properties", "count": {"$sum": 1}}}, - {"$project": {"_id": 1}} + {"$project": {"_id": 1}}, ] result = list(self.featuredb[self.collection].aggregate(pipeline)) @@ -90,13 +90,14 @@ def get_fields(self): # set the field type to 'string'. # by operating without a schema, mongo can query any data type. for i in result: - for key in result[0]['_id'].keys(): - self._fields[key] = {'type': 'string'} + for key in result[0]["_id"].keys(): + self._fields[key] = {"type": "string"} return self._fields - def _get_feature_list(self, filterObj, sortList=[], skip=0, maxitems=1, - skip_geometry=False): + def _get_feature_list( + self, filterObj, sortList=[], skip=0, maxitems=1, skip_geometry=False + ): featurecursor = self.featuredb[self.collection].find(filterObj) if sortList: featurecursor = featurecursor.sort(sortList) @@ -109,18 +110,18 @@ def _get_feature_list(self, filterObj, sortList=[], skip=0, maxitems=1, features = [] for item in featurelist: - feature_id = str(item.pop('_id')) - geometry = item['feature']['geometry'] - props = item['feature']['properties'] - + feature_id = str(item.pop("_id")) + geometry = item["feature"]["geometry"] + props = item["feature"]["properties"] + if skip_geometry: geometry = None feature = { - 'type': 'Feature', - 'id': feature_id, - 'geometry': geometry, - 'properties': props + "type": "Feature", + "id": feature_id, + "geometry": geometry, + "properties": props, } features.append(feature) @@ -128,9 +129,21 @@ def _get_feature_list(self, filterObj, sortList=[], skip=0, maxitems=1, return features @crs_transform - def query(self, offset=0, limit=10, resulttype='results', - bbox=[], datetime_=None, properties=[], sortby=[], - select_properties=[], skip_geometry=False, q=None, filterq=None, **kwargs): + def query( + self, + offset=0, + limit=10, + resulttype="results", + bbox=[], + datetime_=None, + properties=[], + sortby=[], + select_properties=[], + skip_geometry=False, + q=None, + filterq=None, + **kwargs, + ): """ query the provider @@ -140,66 +153,50 @@ def query(self, offset=0, limit=10, resulttype='results', def cql2_to_mongo(node): if node is None: return - + # GeoJson operator if node.__class__.__name__ == "GeometryWithin": field = node.lhs.name geom = node.rhs.geometry - query_body = { - field: { - '$geoWithin': { - '$geometry': geom - } - } - } + query_body = {field: {"$geoWithin": {"$geometry": geom}}} return query_body - + if node.__class__.__name__ == "GeometryIntersects": field = node.lhs.name geom = node.rhs.geometry - return { - field: { - '$geoIntersects': { - '$geometry': geom - } - } - } + return {field: {"$geoIntersects": {"$geometry": geom}}} if node.__class__.__name__ == "DistanceWithin": field = node.lhs.name geom = node.rhs.geometry distance = node.distance - units = node.units # mongo's default units are meters + # units = node.units # mongo's default units are meters # but with CQL we can pass different units # and here we can recalculate them return { field: { - '$near': { - '$geometry': geom, - '$maxDistance': distance, - '$minDistance': 0 + "$near": { + "$geometry": geom, + "$maxDistance": distance, + "$minDistance": 0, } } } # Logical operators if node.__class__.__name__ == "And": - return { - "$and": [ - cql2_to_mongo(node.lhs), - cql2_to_mongo(node.rhs) - ] - } + return {"$and": + [cql2_to_mongo(node.lhs), + cql2_to_mongo(node.rhs)] + } if node.__class__.__name__ == "Or": - return { - "$or": [ - cql2_to_mongo(node.lhs), - cql2_to_mongo(node.rhs) - ] - } + return {"$or": + [cql2_to_mongo(node.lhs), + cql2_to_mongo(node.rhs)] + } # Comparison operators if node.__class__.__name__ == "Equal": @@ -223,38 +220,46 @@ def cql2_to_mongo(node): cql_filters_parsed = cql2_to_mongo(filterq) if cql_filters_parsed is not None: and_filter.append(cql_filters_parsed) - limit = -1 # if there is CQL query return all elements - + limit = -1 # if there is CQL query return all elements + if len(bbox) == 4: x, y, w, h = map(float, bbox) - and_filter.append( - {'geometry': {'$geoWithin': {'$box': [[x, y], [w, h]]}}}) + and_filter.append({ + "geometry": { + "$geoWithin": { + "$box": [[x, y], [w, h]] + } + } + }) - for prop in properties: - and_filter.append({"properties."+prop[0]: {'$eq': prop[1]}}) - - filterobj = {'$and': and_filter} if and_filter else {} + and_filter.append({"properties." + prop[0]: {"$eq": prop[1]}}) - sort_list = [("properties." + sort['property'], - ASCENDING if (sort['order'] == '+') else DESCENDING) - for sort in sortby] + filterobj = {"$and": and_filter} if and_filter else {} - feature_collection = { - 'type': 'FeatureCollection', - 'features': [] - } + sort_list = [ + ( + "properties." + sort["property"], + ASCENDING if (sort["order"] == "+") else DESCENDING, + ) + for sort in sortby + ] - if resulttype == 'hits': + feature_collection = {"type": "FeatureCollection", "features": []} + + if resulttype == "hits": return feature_collection featurelist = self._get_feature_list( - filterobj, sortList=sort_list, skip=offset, maxitems=limit, - skip_geometry=skip_geometry + filterobj, + sortList=sort_list, + skip=offset, + maxitems=limit, + skip_geometry=skip_geometry, ) - feature_collection['features'] = featurelist - feature_collection['numberReturned'] = len(featurelist) + feature_collection["features"] = featurelist + feature_collection["numberReturned"] = len(featurelist) return feature_collection @@ -266,17 +271,16 @@ def get(self, identifier, **kwargs): :param identifier: feature id :returns: dict of single GeoJSON feature """ - featurelist = self._get_feature_list({'_id': ObjectId(identifier)}) + featurelist = self._get_feature_list({"_id": ObjectId(identifier)}) if featurelist: return featurelist[0] else: - err = f'item {identifier} not found' + err = f"item {identifier} not found" LOGGER.error(err) raise ProviderItemNotFoundError(err) def create(self, new_feature): - """Create a new feature - """ + """Create a new feature""" self.featuredb[self.collection].insert_one(new_feature) def update(self, identifier, updated_feature): @@ -285,9 +289,10 @@ def update(self, identifier, updated_feature): :param identifier: feature id :param new_feature: new GeoJSON feature dictionary """ - data = {k: v for k, v in updated_feature.items() if k != 'id'} + data = {k: v for k, v in updated_feature.items() if k != "id"} self.featuredb[self.collection].update_one( - {'_id': ObjectId(identifier)}, {"$set": data}) + {"_id": ObjectId(identifier)}, {"$set": data} + ) def delete(self, identifier): """Deletes an existing feature @@ -295,4 +300,5 @@ def delete(self, identifier): :param identifier: feature id """ self.featuredb[self.collection].delete_one( - {'_id': ObjectId(identifier)}) + {"_id": ObjectId(identifier)} + )