Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@
CheckpointTuple,
get_checkpoint_id,
)
from langgraph.checkpoint.serde.base import SerializerProtocol
from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
from pymongo import ASCENDING, MongoClient, UpdateOne
from pymongo.database import Database as MongoDatabase

Expand Down Expand Up @@ -81,6 +83,7 @@ def __init__(
checkpoint_collection_name: str = "checkpoints",
writes_collection_name: str = "checkpoint_writes",
ttl: Optional[int] = None,
serde: SerializerProtocol | None = None,
**kwargs: Any,
) -> None:
super().__init__()
Expand All @@ -89,6 +92,10 @@ def __init__(
self.checkpoint_collection = self.db[checkpoint_collection_name]
self.writes_collection = self.db[writes_collection_name]
self.ttl = ttl
if serde is not None:
self.serde = serde
else:
self.serde = JsonPlusSerializer()

# Create indexes if not present
if len(self.checkpoint_collection.list_indexes().to_list()) < 2:
Expand Down Expand Up @@ -236,7 +243,7 @@ def get_tuple(self, config: RunnableConfig) -> Optional[CheckpointTuple]:
return CheckpointTuple(
{"configurable": config_values},
checkpoint,
loads_metadata(doc["metadata"]),
loads_metadata(self.serde, doc["metadata"]),
(
{
"configurable": {
Expand Down Expand Up @@ -291,7 +298,7 @@ def list(

if filter:
for key, value in filter.items():
query[f"metadata.{key}"] = dumps_metadata(value)
query[f"metadata.{key}"] = dumps_metadata(self.serde, value)

if before is not None:
query["checkpoint_id"] = {"$lt": before["configurable"]["checkpoint_id"]}
Expand Down Expand Up @@ -325,7 +332,7 @@ def list(
}
},
checkpoint=self.serde.loads_typed((doc["type"], doc["checkpoint"])),
metadata=loads_metadata(doc["metadata"]),
metadata=loads_metadata(self.serde, doc["metadata"]),
parent_config=(
{
"configurable": {
Expand Down Expand Up @@ -381,7 +388,7 @@ def put(
"parent_checkpoint_id": config["configurable"].get("checkpoint_id"),
"type": type_,
"checkpoint": serialized_checkpoint,
"metadata": dumps_metadata(metadata),
"metadata": dumps_metadata(self.serde, metadata),
}
upsert_query = {
"thread_id": thread_id,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,9 @@

from langgraph.checkpoint.base import CheckpointMetadata
from langgraph.checkpoint.serde.base import SerializerProtocol
from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
from pymongo import AsyncMongoClient
from pymongo.driver_info import DriverInfo

serde: SerializerProtocol = JsonPlusSerializer()

DRIVER_METADATA = DriverInfo(
name="Langgraph", version=version("langgraph-checkpoint-mongodb")
)
Expand All @@ -25,7 +22,9 @@ def _append_client_metadata(client: AsyncMongoClient) -> None:
client.append_metadata(DRIVER_METADATA)


def loads_metadata(metadata: dict[str, Any]) -> CheckpointMetadata:
def loads_metadata(
serde: SerializerProtocol, metadata: dict[str, Any]
) -> CheckpointMetadata:
"""Deserialize metadata document

The CheckpointMetadata class itself cannot be stored directly in MongoDB,
Expand All @@ -38,13 +37,14 @@ def loads_metadata(metadata: dict[str, Any]) -> CheckpointMetadata:
if isinstance(metadata, dict):
output = dict()
for key, value in metadata.items():
output[key] = loads_metadata(value)
output[key] = loads_metadata(serde, value)
return output
else:
return serde.loads_typed(metadata)


def dumps_metadata(
serde: SerializerProtocol,
metadata: Union[CheckpointMetadata, Any],
) -> Union[bytes, dict[str, Any]]:
"""Serialize all values in metadata dictionary.
Expand All @@ -54,7 +54,7 @@ def dumps_metadata(
if isinstance(metadata, dict):
output = dict()
for key, value in metadata.items():
output[key] = dumps_metadata(value)
output[key] = dumps_metadata(serde, value)
return output
else:
return serde.dumps_typed(metadata)
56 changes: 56 additions & 0 deletions libs/langgraph-checkpoint-mongodb/tests/unit_tests/test_serde.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
import os
from typing import Any

from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
from pymongo import MongoClient

from langgraph.checkpoint.mongodb import MongoDBSaver

MONGODB_URI = os.environ.get(
"MONGODB_URI", "mongodb://localhost:27017/?directConnection=true"
)
DB_NAME = os.environ.get("DB_NAME", "langgraph-test")
COLLECTION_NAME = "serde_checkpoints"


class CustomSerializer(JsonPlusSerializer):
def __init__(self) -> None:
super().__init__()
self.dumps_called = False
self.loads_called = False

def dumps_typed(self, obj: Any) -> tuple[str, bytes]:
self.dumps_called = True
return super().dumps_typed(obj)

def loads_typed(self, obj: tuple[str, bytes]) -> Any:
self.loads_called = True
return super().loads_typed(obj)


def test_custom_serde(input_data: dict[str, Any]) -> None:
client: MongoClient = MongoClient(MONGODB_URI)
db = client[DB_NAME]
db.drop_collection(COLLECTION_NAME)

custom_serializer = CustomSerializer()

with MongoDBSaver.from_conn_string(
MONGODB_URI, DB_NAME, COLLECTION_NAME, serde=custom_serializer
) as saver:
put_config = saver.put(
input_data["config_1"],
input_data["chkpnt_1"],
input_data["metadata_1"],
{},
)

assert custom_serializer.dumps_called

retrieved_checkpoint_tuple = saver.get_tuple(put_config)

assert custom_serializer.loads_called

assert retrieved_checkpoint_tuple is not None
assert retrieved_checkpoint_tuple.checkpoint == input_data["chkpnt_1"]
assert retrieved_checkpoint_tuple.metadata == input_data["metadata_1"]
Loading