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]
class AzureCosmosDBStore:
 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"}
AzureCosmosDBStore( endpoint: str, key: str, database_name: str, container_name: str, ssl_cert_path: Optional[str] = None, ssl_verify: bool = True)
 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.

def get(self, document_id: str) -> Optional[Dict[str, Any]]:
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.

def upsert(self, document_id: str, data: Dict[str, Any]) -> None:
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.

def delete(self, document_id: str) -> bool:
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.

def exists(self, document_id: str) -> bool:
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.

def list_all(self) -> List[Dict[str, Any]]:
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.

def count(self) -> int:
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.

def clear(self) -> None:
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.

def close(self) -> None:
247    def close(self) -> None:
248        """No-op: azure-cosmos CosmosClient does not require explicit close."""
249        pass

No-op: azure-cosmos CosmosClient does not require explicit close.

class MongoDBStore:
 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"}
MongoDBStore( connection_string: str, database_name: str, collection_name: str, ssl_cert_path: Optional[str] = None)
 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.

def get(self, document_id: str) -> Optional[Dict[str, Any]]:
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.

def upsert(self, document_id: str, data: Dict[str, Any]) -> None:
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.

def delete(self, document_id: str) -> bool:
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.

def exists(self, document_id: str) -> bool:
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.

def list_all(self) -> List[Dict[str, Any]]:
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.

def count(self) -> int:
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.

def clear(self) -> None:
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.

def close(self) -> None:
196    def close(self) -> None:
197        """Close the underlying MongoDB client connection."""
198        self._client.close()

Close the underlying MongoDB client connection.