diff --git a/pygeoapi/api/collection.py b/pygeoapi/api/collection.py index 2cd8b3ea6..584490a62 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 049174f3f..ed064de5d 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([("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,80 +90,176 @@ 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) 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 + + feature = { + "type": "Feature", + "id": feature_id, + "geometry": geometry, + "properties": props, + } + + features.append(feature) - return featurelist + 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): + 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 :returns: dict of 0..n GeoJSON features """ + + 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_}}) + and_filter.append({ + "geometry": { + "$geoWithin": { + "$box": [[x, y], [w, h]] + } + } + }) for prop in properties: - and_filter.append({"properties."+prop[0]: {'$eq': prop[1]}}) + and_filter.append({"properties." + prop[0]: {"$eq": prop[1]}}) - filterobj = {'$and': and_filter} if and_filter else {} + filterobj = {"$and": and_filter} if and_filter else {} - sort_list = [("properties." + sort['property'], - ASCENDING if (sort['order'] == '+') else DESCENDING) - for sort in sortby] + sort_list = [ + ( + "properties." + sort["property"], + ASCENDING if (sort["order"] == "+") else DESCENDING, + ) + for sort in sortby + ] - feature_collection = { - 'type': 'FeatureCollection', - 'features': [] - } + feature_collection = {"type": "FeatureCollection", "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': + 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 @@ -175,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): @@ -194,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 @@ -204,4 +300,5 @@ def delete(self, identifier): :param identifier: feature id """ self.featuredb[self.collection].delete_one( - {'_id': ObjectId(identifier)}) + {"_id": ObjectId(identifier)} + )