gmf_forge_ai_data.data_stores
Data stores module — plain document storage (no embeddings).
Provides lightweight key/value document stores backed by Azure Cosmos DB and MongoDB. These stores are suitable for permanent configuration and metadata storage (orchestrators, agents, user presets) where no vector search or TTL semantics are required.
1""" 2Data stores module — plain document storage (no embeddings). 3 4Provides lightweight key/value document stores backed by Azure Cosmos DB 5and MongoDB. These stores are suitable for permanent configuration and 6metadata storage (orchestrators, agents, user presets) where no vector 7search or TTL semantics are required. 8""" 9 10from .cosmos_store import AzureCosmosDBStore 11from .mongo_store import MongoDBStore 12 13__all__ = [ 14 "AzureCosmosDBStore", 15 "MongoDBStore", 16]
29class AzureCosmosDBStore: 30 """ 31 Azure Cosmos DB NoSQL document store (no embeddings). 32 33 Stores arbitrary dicts as Cosmos DB items. ``document_id`` becomes the 34 item ``id`` and partition key; all caller fields are stored at the top 35 level. No embedding, content, timestamp, or document_data wrapper fields 36 are ever written. 37 38 The database and container are provisioned automatically with 39 ``/id`` as the partition key path. 40 41 Usage:: 42 43 from gmf_forge_ai_data.data_stores import AzureCosmosDBStore 44 45 store = AzureCosmosDBStore( 46 endpoint="https://your-account.documents.azure.com:443/", 47 key="your-key", 48 database_name="portal_db", 49 container_name="portal-orchestrators", 50 ) 51 store.upsert("orch-1", {"name": "Main", "status": "healthy"}) 52 doc = store.get("orch-1") # {"name": "Main", "status": "healthy"} 53 """ 54 55 #: Cosmos DB system fields stripped from all read results. 56 _SYSTEM_FIELDS: frozenset = frozenset( 57 {"_rid", "_self", "_etag", "_attachments", "_ts"} 58 ) 59 60 def __init__( 61 self, 62 endpoint: str, 63 key: str, 64 database_name: str, 65 container_name: str, 66 ssl_cert_path: Optional[str] = None, 67 ssl_verify: bool = True, 68 ) -> None: 69 """ 70 Initialise the Cosmos DB document store. 71 72 Args: 73 endpoint: Cosmos DB account endpoint 74 (e.g. ``https://your-account.documents.azure.com:443/``). 75 key: Cosmos DB account key or resource token. 76 database_name: Name of the Cosmos DB database. 77 container_name: Name of the container. 78 ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS 79 verification. Useful in corporate environments with custom 80 certificate authorities. 81 ssl_verify: Set to ``False`` to disable TLS certificate 82 verification. Use only for local emulators or dev environments. 83 """ 84 try: 85 from azure.cosmos import CosmosClient # noqa: F401 86 except ImportError as exc: 87 raise ImportError( 88 "azure-cosmos is required for AzureCosmosDBStore. " 89 "Install it with: pip install azure-cosmos" 90 ) from exc 91 92 self._container_name = container_name 93 94 connection_kwargs: Dict[str, Any] = {} 95 if ssl_cert_path: 96 import ssl as _ssl 97 _ssl.create_default_context(cafile=ssl_cert_path) 98 connection_kwargs["connection_verify"] = ssl_cert_path 99 elif not ssl_verify: 100 connection_kwargs["connection_verify"] = False 101 102 from azure.cosmos import CosmosClient, PartitionKey 103 from azure.cosmos.documents import ConnectionPolicy 104 105 # When ssl_verify=False (local emulator) the emulator's account-read response 106 # contains 127.0.0.1 as its partition endpoints. From inside a Docker container 107 # 127.0.0.1 resolves to the container itself, so all subsequent SDK requests 108 # fail with "Connection refused". Disabling endpoint discovery keeps every 109 # request on the originally supplied endpoint URL. 110 if not ssl_verify: 111 _policy = ConnectionPolicy() 112 _policy.EnableEndpointDiscovery = False 113 connection_kwargs["connection_policy"] = _policy 114 115 self._client = CosmosClient(endpoint, credential=key, **connection_kwargs) 116 self._database = self._client.create_database_if_not_exists(id=database_name) 117 self._container = self._database.create_container_if_not_exists( 118 id=container_name, 119 partition_key=PartitionKey(path="/id"), 120 ) 121 122 # ------------------------------------------------------------------ 123 # Internal helpers 124 # ------------------------------------------------------------------ 125 126 def _clean(self, item: Dict[str, Any]) -> Dict[str, Any]: 127 """Strip Cosmos system fields and the ``id`` field from a raw item.""" 128 return { 129 k: v 130 for k, v in item.items() 131 if k != "id" and k not in self._SYSTEM_FIELDS 132 } 133 134 # ------------------------------------------------------------------ 135 # CRUD interface 136 # ------------------------------------------------------------------ 137 138 def get(self, document_id: str) -> Optional[Dict[str, Any]]: 139 """ 140 Retrieve a document by ID. 141 142 Args: 143 document_id: The document's unique identifier. 144 145 Returns: 146 Dict of document fields (system fields and ``id`` stripped), 147 or ``None`` if the document does not exist. 148 """ 149 try: 150 from azure.cosmos.exceptions import CosmosResourceNotFoundError 151 item = self._container.read_item( 152 item=document_id, partition_key=document_id 153 ) 154 return self._clean(item) 155 except Exception as exc: 156 try: 157 from azure.cosmos.exceptions import CosmosResourceNotFoundError as _NotFound 158 except ImportError: 159 raise 160 if isinstance(exc, _NotFound): 161 return None 162 raise 163 164 def upsert(self, document_id: str, data: Dict[str, Any]) -> None: 165 """ 166 Insert or replace a document. 167 168 The ``document_id`` is stored as the Cosmos ``id`` and partition key. 169 All fields in *data* are stored at the top level alongside ``id``. 170 171 Args: 172 document_id: Unique identifier for the document. 173 data: Payload to store. Must not contain an ``"id"`` key. 174 """ 175 self._container.upsert_item({"id": document_id, **data}) 176 177 def delete(self, document_id: str) -> bool: 178 """ 179 Delete a document by ID. 180 181 Args: 182 document_id: The document's unique identifier. 183 184 Returns: 185 ``True`` if the document was deleted, ``False`` if it did not exist. 186 """ 187 try: 188 self._container.delete_item(item=document_id, partition_key=document_id) 189 return True 190 except Exception as exc: 191 try: 192 from azure.cosmos.exceptions import CosmosResourceNotFoundError as _NotFound 193 except ImportError: 194 raise 195 if isinstance(exc, _NotFound): 196 return False 197 raise 198 199 def exists(self, document_id: str) -> bool: 200 """Return ``True`` if a document with the given ID exists.""" 201 return self.get(document_id) is not None 202 203 def list_all(self) -> List[Dict[str, Any]]: 204 """ 205 Return all documents in the container. 206 207 Returns: 208 List of dicts, each with system fields and ``id`` stripped. 209 """ 210 items = list(self._container.query_items( 211 query="SELECT * FROM c", 212 enable_cross_partition_query=True, 213 )) 214 return [self._clean(item) for item in items] 215 216 def count(self) -> int: 217 """Return the total number of documents in the container.""" 218 items = list(self._container.query_items( 219 query="SELECT VALUE COUNT(1) FROM c", 220 enable_cross_partition_query=True, 221 )) 222 return items[0] if items else 0 223 224 def clear(self) -> None: 225 """ 226 Remove **all** documents from the container. 227 228 Warning: 229 This operation is irreversible. 230 """ 231 items = list(self._container.query_items( 232 query="SELECT c.id FROM c", 233 enable_cross_partition_query=True, 234 )) 235 for item in items: 236 try: 237 self._container.delete_item( 238 item=item["id"], partition_key=item["id"] 239 ) 240 except Exception: 241 pass 242 logger.warning( 243 "Cleared all documents from Cosmos DB container '%s'", 244 self._container_name, 245 ) 246 247 def close(self) -> None: 248 """No-op: azure-cosmos CosmosClient does not require explicit close.""" 249 pass
Azure Cosmos DB NoSQL document store (no embeddings).
Stores arbitrary dicts as Cosmos DB items. document_id becomes the
item id and partition key; all caller fields are stored at the top
level. No embedding, content, timestamp, or document_data wrapper fields
are ever written.
The database and container are provisioned automatically with
/id as the partition key path.
Usage::
from gmf_forge_ai_data.data_stores import AzureCosmosDBStore
store = AzureCosmosDBStore(
endpoint="https://your-account.documents.azure.com:443/",
key="your-key",
database_name="portal_db",
container_name="portal-orchestrators",
)
store.upsert("orch-1", {"name": "Main", "status": "healthy"})
doc = store.get("orch-1") # {"name": "Main", "status": "healthy"}
60 def __init__( 61 self, 62 endpoint: str, 63 key: str, 64 database_name: str, 65 container_name: str, 66 ssl_cert_path: Optional[str] = None, 67 ssl_verify: bool = True, 68 ) -> None: 69 """ 70 Initialise the Cosmos DB document store. 71 72 Args: 73 endpoint: Cosmos DB account endpoint 74 (e.g. ``https://your-account.documents.azure.com:443/``). 75 key: Cosmos DB account key or resource token. 76 database_name: Name of the Cosmos DB database. 77 container_name: Name of the container. 78 ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS 79 verification. Useful in corporate environments with custom 80 certificate authorities. 81 ssl_verify: Set to ``False`` to disable TLS certificate 82 verification. Use only for local emulators or dev environments. 83 """ 84 try: 85 from azure.cosmos import CosmosClient # noqa: F401 86 except ImportError as exc: 87 raise ImportError( 88 "azure-cosmos is required for AzureCosmosDBStore. " 89 "Install it with: pip install azure-cosmos" 90 ) from exc 91 92 self._container_name = container_name 93 94 connection_kwargs: Dict[str, Any] = {} 95 if ssl_cert_path: 96 import ssl as _ssl 97 _ssl.create_default_context(cafile=ssl_cert_path) 98 connection_kwargs["connection_verify"] = ssl_cert_path 99 elif not ssl_verify: 100 connection_kwargs["connection_verify"] = False 101 102 from azure.cosmos import CosmosClient, PartitionKey 103 from azure.cosmos.documents import ConnectionPolicy 104 105 # When ssl_verify=False (local emulator) the emulator's account-read response 106 # contains 127.0.0.1 as its partition endpoints. From inside a Docker container 107 # 127.0.0.1 resolves to the container itself, so all subsequent SDK requests 108 # fail with "Connection refused". Disabling endpoint discovery keeps every 109 # request on the originally supplied endpoint URL. 110 if not ssl_verify: 111 _policy = ConnectionPolicy() 112 _policy.EnableEndpointDiscovery = False 113 connection_kwargs["connection_policy"] = _policy 114 115 self._client = CosmosClient(endpoint, credential=key, **connection_kwargs) 116 self._database = self._client.create_database_if_not_exists(id=database_name) 117 self._container = self._database.create_container_if_not_exists( 118 id=container_name, 119 partition_key=PartitionKey(path="/id"), 120 )
Initialise the Cosmos DB document store.
Args:
endpoint: Cosmos DB account endpoint
(e.g. https://your-account.documents.azure.com:443/).
key: Cosmos DB account key or resource token.
database_name: Name of the Cosmos DB database.
container_name: Name of the container.
ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS
verification. Useful in corporate environments with custom
certificate authorities.
ssl_verify: Set to False to disable TLS certificate
verification. Use only for local emulators or dev environments.
138 def get(self, document_id: str) -> Optional[Dict[str, Any]]: 139 """ 140 Retrieve a document by ID. 141 142 Args: 143 document_id: The document's unique identifier. 144 145 Returns: 146 Dict of document fields (system fields and ``id`` stripped), 147 or ``None`` if the document does not exist. 148 """ 149 try: 150 from azure.cosmos.exceptions import CosmosResourceNotFoundError 151 item = self._container.read_item( 152 item=document_id, partition_key=document_id 153 ) 154 return self._clean(item) 155 except Exception as exc: 156 try: 157 from azure.cosmos.exceptions import CosmosResourceNotFoundError as _NotFound 158 except ImportError: 159 raise 160 if isinstance(exc, _NotFound): 161 return None 162 raise
Retrieve a document by ID.
Args: document_id: The document's unique identifier.
Returns:
Dict of document fields (system fields and id stripped),
or None if the document does not exist.
164 def upsert(self, document_id: str, data: Dict[str, Any]) -> None: 165 """ 166 Insert or replace a document. 167 168 The ``document_id`` is stored as the Cosmos ``id`` and partition key. 169 All fields in *data* are stored at the top level alongside ``id``. 170 171 Args: 172 document_id: Unique identifier for the document. 173 data: Payload to store. Must not contain an ``"id"`` key. 174 """ 175 self._container.upsert_item({"id": document_id, **data})
Insert or replace a document.
The document_id is stored as the Cosmos id and partition key.
All fields in data are stored at the top level alongside id.
Args:
document_id: Unique identifier for the document.
data: Payload to store. Must not contain an "id" key.
177 def delete(self, document_id: str) -> bool: 178 """ 179 Delete a document by ID. 180 181 Args: 182 document_id: The document's unique identifier. 183 184 Returns: 185 ``True`` if the document was deleted, ``False`` if it did not exist. 186 """ 187 try: 188 self._container.delete_item(item=document_id, partition_key=document_id) 189 return True 190 except Exception as exc: 191 try: 192 from azure.cosmos.exceptions import CosmosResourceNotFoundError as _NotFound 193 except ImportError: 194 raise 195 if isinstance(exc, _NotFound): 196 return False 197 raise
Delete a document by ID.
Args: document_id: The document's unique identifier.
Returns:
True if the document was deleted, False if it did not exist.
199 def exists(self, document_id: str) -> bool: 200 """Return ``True`` if a document with the given ID exists.""" 201 return self.get(document_id) is not None
Return True if a document with the given ID exists.
203 def list_all(self) -> List[Dict[str, Any]]: 204 """ 205 Return all documents in the container. 206 207 Returns: 208 List of dicts, each with system fields and ``id`` stripped. 209 """ 210 items = list(self._container.query_items( 211 query="SELECT * FROM c", 212 enable_cross_partition_query=True, 213 )) 214 return [self._clean(item) for item in items]
Return all documents in the container.
Returns:
List of dicts, each with system fields and id stripped.
216 def count(self) -> int: 217 """Return the total number of documents in the container.""" 218 items = list(self._container.query_items( 219 query="SELECT VALUE COUNT(1) FROM c", 220 enable_cross_partition_query=True, 221 )) 222 return items[0] if items else 0
Return the total number of documents in the container.
224 def clear(self) -> None: 225 """ 226 Remove **all** documents from the container. 227 228 Warning: 229 This operation is irreversible. 230 """ 231 items = list(self._container.query_items( 232 query="SELECT c.id FROM c", 233 enable_cross_partition_query=True, 234 )) 235 for item in items: 236 try: 237 self._container.delete_item( 238 item=item["id"], partition_key=item["id"] 239 ) 240 except Exception: 241 pass 242 logger.warning( 243 "Cleared all documents from Cosmos DB container '%s'", 244 self._container_name, 245 )
Remove all documents from the container.
Warning: This operation is irreversible.
37class MongoDBStore: 38 """ 39 MongoDB document store (no embeddings). 40 41 Stores arbitrary dicts in a MongoDB collection using ``_doc_id`` as 42 the lookup key. All caller fields are stored at the top level. No 43 embedding, content, or vector-search fields are ever written. 44 45 A unique index on ``_doc_id`` is created automatically on first use. 46 47 Usage:: 48 49 from gmf_forge_ai_data.data_stores import MongoDBStore 50 51 store = MongoDBStore( 52 connection_string="mongodb+srv://...", 53 database_name="portal_db", 54 collection_name="portal-orchestrators", 55 ) 56 store.upsert("orch-1", {"name": "Main", "status": "healthy"}) 57 doc = store.get("orch-1") # {"name": "Main", "status": "healthy"} 58 """ 59 60 #: Internal fields stripped from all read results. 61 _INTERNAL_FIELDS: frozenset = frozenset({"_id", _DOC_ID_FIELD}) 62 63 def __init__( 64 self, 65 connection_string: str, 66 database_name: str, 67 collection_name: str, 68 ssl_cert_path: Optional[str] = None, 69 ) -> None: 70 """ 71 Initialise the MongoDB document store. 72 73 A unique index on the ``_doc_id`` field is created automatically if it 74 does not already exist. 75 76 Args: 77 connection_string: MongoDB connection string 78 (e.g. ``mongodb+srv://user:pass@cluster.mongodb.net/``). 79 database_name: Name of the MongoDB database. 80 collection_name: Name of the collection. 81 ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS 82 verification. Useful in corporate environments with custom 83 certificate authorities. 84 """ 85 try: 86 import pymongo # noqa: F401 87 except ImportError as exc: 88 raise ImportError( 89 "pymongo is required for MongoDBStore. " 90 "Install it with: pip install pymongo" 91 ) from exc 92 93 self._collection_name = collection_name 94 95 client_kwargs: Dict[str, Any] = {} 96 if ssl_cert_path: 97 client_kwargs["tlsCAFile"] = ssl_cert_path 98 99 import pymongo 100 101 self._client: pymongo.MongoClient = pymongo.MongoClient( 102 connection_string, **client_kwargs 103 ) 104 self._db = self._client[database_name] 105 self._collection = self._db[collection_name] 106 self._collection.create_index(_DOC_ID_FIELD, unique=True) 107 108 # ------------------------------------------------------------------ 109 # Internal helpers 110 # ------------------------------------------------------------------ 111 112 def _clean(self, doc: Dict[str, Any]) -> Dict[str, Any]: 113 """Strip ``_id`` and ``_doc_id`` from a raw MongoDB document.""" 114 return {k: v for k, v in doc.items() if k not in self._INTERNAL_FIELDS} 115 116 # ------------------------------------------------------------------ 117 # CRUD interface 118 # ------------------------------------------------------------------ 119 120 def get(self, document_id: str) -> Optional[Dict[str, Any]]: 121 """ 122 Retrieve a document by ID. 123 124 Args: 125 document_id: The document's unique identifier. 126 127 Returns: 128 Dict of document fields (``_id`` and ``_doc_id`` stripped), 129 or ``None`` if the document does not exist. 130 """ 131 doc = self._collection.find_one({_DOC_ID_FIELD: document_id}) 132 if doc is None: 133 return None 134 return self._clean(doc) 135 136 def upsert(self, document_id: str, data: Dict[str, Any]) -> None: 137 """ 138 Insert or replace a document. 139 140 Args: 141 document_id: Unique identifier for the document. 142 data: Payload to store. Must not contain ``"_doc_id"`` or 143 ``"_id"`` keys. 144 """ 145 self._collection.update_one( 146 {_DOC_ID_FIELD: document_id}, 147 {"$set": {_DOC_ID_FIELD: document_id, **data}}, 148 upsert=True, 149 ) 150 151 def delete(self, document_id: str) -> bool: 152 """ 153 Delete a document by ID. 154 155 Args: 156 document_id: The document's unique identifier. 157 158 Returns: 159 ``True`` if the document was deleted, ``False`` if it did not exist. 160 """ 161 result = self._collection.delete_one({_DOC_ID_FIELD: document_id}) 162 return result.deleted_count > 0 163 164 def exists(self, document_id: str) -> bool: 165 """Return ``True`` if a document with the given ID exists.""" 166 return ( 167 self._collection.count_documents({_DOC_ID_FIELD: document_id}, limit=1) > 0 168 ) 169 170 def list_all(self) -> List[Dict[str, Any]]: 171 """ 172 Return all documents in the collection. 173 174 Returns: 175 List of dicts, each with ``_id`` and ``_doc_id`` stripped. 176 """ 177 return [self._clean(doc) for doc in self._collection.find({})] 178 179 def count(self) -> int: 180 """Return the total number of documents in the collection.""" 181 return self._collection.count_documents({}) 182 183 def clear(self) -> None: 184 """ 185 Remove **all** documents from the collection. 186 187 Warning: 188 This operation is irreversible. 189 """ 190 self._collection.delete_many({}) 191 logger.warning( 192 "Cleared all documents from MongoDB collection '%s'", 193 self._collection_name, 194 ) 195 196 def close(self) -> None: 197 """Close the underlying MongoDB client connection.""" 198 self._client.close()
MongoDB document store (no embeddings).
Stores arbitrary dicts in a MongoDB collection using _doc_id as
the lookup key. All caller fields are stored at the top level. No
embedding, content, or vector-search fields are ever written.
A unique index on _doc_id is created automatically on first use.
Usage::
from gmf_forge_ai_data.data_stores import MongoDBStore
store = MongoDBStore(
connection_string="mongodb+srv://...",
database_name="portal_db",
collection_name="portal-orchestrators",
)
store.upsert("orch-1", {"name": "Main", "status": "healthy"})
doc = store.get("orch-1") # {"name": "Main", "status": "healthy"}
63 def __init__( 64 self, 65 connection_string: str, 66 database_name: str, 67 collection_name: str, 68 ssl_cert_path: Optional[str] = None, 69 ) -> None: 70 """ 71 Initialise the MongoDB document store. 72 73 A unique index on the ``_doc_id`` field is created automatically if it 74 does not already exist. 75 76 Args: 77 connection_string: MongoDB connection string 78 (e.g. ``mongodb+srv://user:pass@cluster.mongodb.net/``). 79 database_name: Name of the MongoDB database. 80 collection_name: Name of the collection. 81 ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS 82 verification. Useful in corporate environments with custom 83 certificate authorities. 84 """ 85 try: 86 import pymongo # noqa: F401 87 except ImportError as exc: 88 raise ImportError( 89 "pymongo is required for MongoDBStore. " 90 "Install it with: pip install pymongo" 91 ) from exc 92 93 self._collection_name = collection_name 94 95 client_kwargs: Dict[str, Any] = {} 96 if ssl_cert_path: 97 client_kwargs["tlsCAFile"] = ssl_cert_path 98 99 import pymongo 100 101 self._client: pymongo.MongoClient = pymongo.MongoClient( 102 connection_string, **client_kwargs 103 ) 104 self._db = self._client[database_name] 105 self._collection = self._db[collection_name] 106 self._collection.create_index(_DOC_ID_FIELD, unique=True)
Initialise the MongoDB document store.
A unique index on the _doc_id field is created automatically if it
does not already exist.
Args:
connection_string: MongoDB connection string
(e.g. mongodb+srv://user:pass@cluster.mongodb.net/).
database_name: Name of the MongoDB database.
collection_name: Name of the collection.
ssl_cert_path: Path to a CA certificate bundle (PEM) for TLS
verification. Useful in corporate environments with custom
certificate authorities.
120 def get(self, document_id: str) -> Optional[Dict[str, Any]]: 121 """ 122 Retrieve a document by ID. 123 124 Args: 125 document_id: The document's unique identifier. 126 127 Returns: 128 Dict of document fields (``_id`` and ``_doc_id`` stripped), 129 or ``None`` if the document does not exist. 130 """ 131 doc = self._collection.find_one({_DOC_ID_FIELD: document_id}) 132 if doc is None: 133 return None 134 return self._clean(doc)
Retrieve a document by ID.
Args: document_id: The document's unique identifier.
Returns:
Dict of document fields (_id and _doc_id stripped),
or None if the document does not exist.
136 def upsert(self, document_id: str, data: Dict[str, Any]) -> None: 137 """ 138 Insert or replace a document. 139 140 Args: 141 document_id: Unique identifier for the document. 142 data: Payload to store. Must not contain ``"_doc_id"`` or 143 ``"_id"`` keys. 144 """ 145 self._collection.update_one( 146 {_DOC_ID_FIELD: document_id}, 147 {"$set": {_DOC_ID_FIELD: document_id, **data}}, 148 upsert=True, 149 )
Insert or replace a document.
Args:
document_id: Unique identifier for the document.
data: Payload to store. Must not contain "_doc_id" or
"_id" keys.
151 def delete(self, document_id: str) -> bool: 152 """ 153 Delete a document by ID. 154 155 Args: 156 document_id: The document's unique identifier. 157 158 Returns: 159 ``True`` if the document was deleted, ``False`` if it did not exist. 160 """ 161 result = self._collection.delete_one({_DOC_ID_FIELD: document_id}) 162 return result.deleted_count > 0
Delete a document by ID.
Args: document_id: The document's unique identifier.
Returns:
True if the document was deleted, False if it did not exist.
164 def exists(self, document_id: str) -> bool: 165 """Return ``True`` if a document with the given ID exists.""" 166 return ( 167 self._collection.count_documents({_DOC_ID_FIELD: document_id}, limit=1) > 0 168 )
Return True if a document with the given ID exists.
170 def list_all(self) -> List[Dict[str, Any]]: 171 """ 172 Return all documents in the collection. 173 174 Returns: 175 List of dicts, each with ``_id`` and ``_doc_id`` stripped. 176 """ 177 return [self._clean(doc) for doc in self._collection.find({})]
Return all documents in the collection.
Returns:
List of dicts, each with _id and _doc_id stripped.
179 def count(self) -> int: 180 """Return the total number of documents in the collection.""" 181 return self._collection.count_documents({})
Return the total number of documents in the collection.
183 def clear(self) -> None: 184 """ 185 Remove **all** documents from the collection. 186 187 Warning: 188 This operation is irreversible. 189 """ 190 self._collection.delete_many({}) 191 logger.warning( 192 "Cleared all documents from MongoDB collection '%s'", 193 self._collection_name, 194 )
Remove all documents from the collection.
Warning: This operation is irreversible.