|
1 | 1 | import os |
2 | | -from typing import Any |
| 2 | +from collections.abc import AsyncGenerator |
| 3 | +from typing import Any, Union |
3 | 4 |
|
4 | 5 | import pytest |
| 6 | +import pytest_asyncio |
5 | 7 | from bson.errors import InvalidDocument |
6 | | -from pymongo import AsyncMongoClient |
| 8 | +from pymongo import AsyncMongoClient, MongoClient |
7 | 9 |
|
8 | | -from langgraph.checkpoint.mongodb.aio import AsyncMongoDBSaver |
| 10 | +from langgraph.checkpoint.mongodb import AsyncMongoDBSaver, MongoDBSaver |
9 | 11 |
|
10 | | -MONGODB_URI = os.environ.get("MONGODB_URI", "mongodb://localhost:27017") |
| 12 | +MONGODB_URI = os.environ.get( |
| 13 | + "MONGODB_URI", "mongodb://localhost:27017/?directConnection=true" |
| 14 | +) |
11 | 15 | DB_NAME = os.environ.get("DB_NAME", "langgraph-test") |
12 | 16 | COLLECTION_NAME = "sync_checkpoints_aio" |
13 | 17 |
|
14 | 18 |
|
15 | | -async def test_asearch(input_data: dict[str, Any]) -> None: |
16 | | - # Clear collections if they exist |
17 | | - client: AsyncMongoClient = AsyncMongoClient(MONGODB_URI) |
18 | | - db = client[DB_NAME] |
19 | | - |
20 | | - for clxn in await db.list_collection_names(): |
21 | | - await db.drop_collection(clxn) |
22 | | - |
23 | | - async with AsyncMongoDBSaver.from_conn_string( |
24 | | - MONGODB_URI, DB_NAME, COLLECTION_NAME |
25 | | - ) as saver: |
26 | | - # save checkpoints |
27 | | - await saver.aput( |
28 | | - input_data["config_1"], |
29 | | - input_data["chkpnt_1"], |
30 | | - input_data["metadata_1"], |
31 | | - {}, |
32 | | - ) |
33 | | - await saver.aput( |
34 | | - input_data["config_2"], |
35 | | - input_data["chkpnt_2"], |
36 | | - input_data["metadata_2"], |
37 | | - {}, |
38 | | - ) |
39 | | - await saver.aput( |
40 | | - input_data["config_3"], |
41 | | - input_data["chkpnt_3"], |
42 | | - input_data["metadata_3"], |
43 | | - {}, |
44 | | - ) |
45 | | - |
46 | | - # call method / assertions |
47 | | - query_1 = {"source": "input"} # search by 1 key |
48 | | - query_2 = { |
49 | | - "step": 1, |
50 | | - "writes": {"foo": "bar"}, |
51 | | - } # search by multiple keys |
52 | | - query_3: dict[str, Any] = {} # search by no keys, return all checkpoints |
53 | | - query_4 = {"source": "update", "step": 1} # no match |
54 | | - |
55 | | - search_results_1 = [c async for c in saver.alist(None, filter=query_1)] |
56 | | - assert len(search_results_1) == 1 |
57 | | - assert search_results_1[0].metadata == input_data["metadata_1"] |
58 | | - |
59 | | - search_results_2 = [c async for c in saver.alist(None, filter=query_2)] |
60 | | - assert len(search_results_2) == 1 |
61 | | - assert search_results_2[0].metadata == input_data["metadata_2"] |
62 | | - |
63 | | - search_results_3 = [c async for c in saver.alist(None, filter=query_3)] |
64 | | - assert len(search_results_3) == 3 |
65 | | - |
66 | | - search_results_4 = [c async for c in saver.alist(None, filter=query_4)] |
67 | | - assert len(search_results_4) == 0 |
68 | | - |
69 | | - # search by config (defaults to checkpoints across all namespaces) |
70 | | - search_results_5 = [ |
71 | | - c async for c in saver.alist({"configurable": {"thread_id": "thread-2"}}) |
72 | | - ] |
73 | | - assert len(search_results_5) == 2 |
74 | | - assert { |
75 | | - search_results_5[0].config["configurable"]["checkpoint_ns"], |
76 | | - search_results_5[1].config["configurable"]["checkpoint_ns"], |
77 | | - } == {"", "inner"} |
78 | | - |
79 | | - |
80 | | -async def test_null_chars(input_data: dict[str, Any]) -> None: |
| 19 | +@pytest_asyncio.fixture(params=["run_in_executor", "aio"]) |
| 20 | +async def async_saver(request: pytest.FixtureRequest) -> AsyncGenerator: |
| 21 | + if request.param == "aio": |
| 22 | + # Use async client and checkpointer |
| 23 | + aclient: AsyncMongoClient = AsyncMongoClient(MONGODB_URI) |
| 24 | + adb = aclient[DB_NAME] |
| 25 | + for clxn in await adb.list_collection_names(): |
| 26 | + await adb.drop_collection(clxn) |
| 27 | + async with AsyncMongoDBSaver.from_conn_string( |
| 28 | + MONGODB_URI, DB_NAME, COLLECTION_NAME |
| 29 | + ) as checkpointer: |
| 30 | + yield checkpointer |
| 31 | + await aclient.close() |
| 32 | + else: |
| 33 | + # Use sync client and checkpointer with async methods run in executor |
| 34 | + client: MongoClient = MongoClient(MONGODB_URI) |
| 35 | + db = client[DB_NAME] |
| 36 | + for clxn in db.list_collection_names(): |
| 37 | + db.drop_collection(clxn) |
| 38 | + with MongoDBSaver.from_conn_string( |
| 39 | + MONGODB_URI, DB_NAME, COLLECTION_NAME |
| 40 | + ) as checkpointer: |
| 41 | + yield checkpointer |
| 42 | + client.close() |
| 43 | + |
| 44 | + |
| 45 | +@pytest.mark.asyncio |
| 46 | +async def test_asearch( |
| 47 | + input_data: dict[str, Any], async_saver: Union[AsyncMongoDBSaver, MongoDBSaver] |
| 48 | +) -> None: |
| 49 | + # save checkpoints |
| 50 | + await async_saver.aput( |
| 51 | + input_data["config_1"], |
| 52 | + input_data["chkpnt_1"], |
| 53 | + input_data["metadata_1"], |
| 54 | + {}, |
| 55 | + ) |
| 56 | + await async_saver.aput( |
| 57 | + input_data["config_2"], |
| 58 | + input_data["chkpnt_2"], |
| 59 | + input_data["metadata_2"], |
| 60 | + {}, |
| 61 | + ) |
| 62 | + await async_saver.aput( |
| 63 | + input_data["config_3"], |
| 64 | + input_data["chkpnt_3"], |
| 65 | + input_data["metadata_3"], |
| 66 | + {}, |
| 67 | + ) |
| 68 | + |
| 69 | + # call method / assertions |
| 70 | + query_1 = {"source": "input"} # search by 1 key |
| 71 | + query_2 = { |
| 72 | + "step": 1, |
| 73 | + "writes": {"foo": "bar"}, |
| 74 | + } # search by multiple keys |
| 75 | + query_3: dict[str, Any] = {} # search by no keys, return all checkpoints |
| 76 | + query_4 = {"source": "update", "step": 1} # no match |
| 77 | + |
| 78 | + search_results_1 = [c async for c in async_saver.alist(None, filter=query_1)] |
| 79 | + assert len(search_results_1) == 1 |
| 80 | + assert search_results_1[0].metadata == input_data["metadata_1"] |
| 81 | + |
| 82 | + search_results_2 = [c async for c in async_saver.alist(None, filter=query_2)] |
| 83 | + assert len(search_results_2) == 1 |
| 84 | + assert search_results_2[0].metadata == input_data["metadata_2"] |
| 85 | + |
| 86 | + search_results_3 = [c async for c in async_saver.alist(None, filter=query_3)] |
| 87 | + assert len(search_results_3) == 3 |
| 88 | + |
| 89 | + search_results_4 = [c async for c in async_saver.alist(None, filter=query_4)] |
| 90 | + assert len(search_results_4) == 0 |
| 91 | + |
| 92 | + # search by config (defaults to checkpoints across all namespaces) |
| 93 | + search_results_5 = [ |
| 94 | + c async for c in async_saver.alist({"configurable": {"thread_id": "thread-2"}}) |
| 95 | + ] |
| 96 | + assert len(search_results_5) == 2 |
| 97 | + assert { |
| 98 | + search_results_5[0].config["configurable"]["checkpoint_ns"], |
| 99 | + search_results_5[1].config["configurable"]["checkpoint_ns"], |
| 100 | + } == {"", "inner"} |
| 101 | + |
| 102 | + |
| 103 | +@pytest.mark.asyncio |
| 104 | +async def test_null_chars( |
| 105 | + input_data: dict[str, Any], async_saver: Union[AsyncMongoDBSaver, MongoDBSaver] |
| 106 | +) -> None: |
81 | 107 | """In MongoDB string *values* can be any valid UTF-8 including nulls. |
82 | 108 | *Field names*, however, cannot contain nulls characters.""" |
83 | | - async with AsyncMongoDBSaver.from_conn_string( |
84 | | - MONGODB_URI, DB_NAME, COLLECTION_NAME |
85 | | - ) as saver: |
86 | | - null_str = "\x00abc" # string containing null character |
87 | 109 |
|
88 | | - # 1. null string in field *value* |
89 | | - null_value_cfg = await saver.aput( |
| 110 | + null_str = "\x00abc" # string containing null character |
| 111 | + |
| 112 | + # 1. null string in field *value* |
| 113 | + null_value_cfg = await async_saver.aput( |
| 114 | + input_data["config_1"], |
| 115 | + input_data["chkpnt_1"], |
| 116 | + {"my_key": null_str}, |
| 117 | + {}, |
| 118 | + ) |
| 119 | + null_tuple = await async_saver.aget_tuple(null_value_cfg) |
| 120 | + assert null_tuple.metadata["my_key"] == null_str # type: ignore |
| 121 | + cps = [c async for c in async_saver.alist(None, filter={"my_key": null_str})] |
| 122 | + assert cps[0].metadata["my_key"] == null_str |
| 123 | + |
| 124 | + # 2. null string in field *name* |
| 125 | + with pytest.raises(InvalidDocument): |
| 126 | + await async_saver.aput( |
90 | 127 | input_data["config_1"], |
91 | 128 | input_data["chkpnt_1"], |
92 | | - {"my_key": null_str}, |
| 129 | + {null_str: "my_value"}, # type: ignore |
93 | 130 | {}, |
94 | 131 | ) |
95 | | - null_tuple = await saver.aget_tuple(null_value_cfg) |
96 | | - assert null_tuple.metadata["my_key"] == null_str # type: ignore |
97 | | - cps = [c async for c in saver.alist(None, filter={"my_key": null_str})] |
98 | | - assert cps[0].metadata["my_key"] == null_str |
99 | | - |
100 | | - # 2. null string in field *name* |
101 | | - with pytest.raises(InvalidDocument): |
102 | | - await saver.aput( |
103 | | - input_data["config_1"], |
104 | | - input_data["chkpnt_1"], |
105 | | - {null_str: "my_value"}, # type: ignore |
106 | | - {}, |
107 | | - ) |
|
0 commit comments