diff --git a/libs/langgraph/poetry.lock b/libs/langgraph/poetry.lock index c9bf1c7fe..93d04e327 100644 --- a/libs/langgraph/poetry.lock +++ b/libs/langgraph/poetry.lock @@ -1255,7 +1255,7 @@ typing-extensions = ">=4.7" [[package]] name = "langgraph-checkpoint" -version = "1.0.8" +version = "1.0.9" description = "Library with base interfaces for LangGraph checkpoint savers." optional = false python-versions = "^3.9.0,<4.0" @@ -2956,6 +2956,50 @@ h2 = ["h2 (>=4,<5)"] socks = ["pysocks (>=1.5.6,!=1.5.7,<2.0)"] zstd = ["zstandard (>=0.18.0)"] +[[package]] +name = "uvloop" +version = "0.20.0" +description = "Fast implementation of asyncio event loop on top of libuv" +optional = false +python-versions = ">=3.8.0" +files = [ + {file = "uvloop-0.20.0-cp310-cp310-macosx_10_9_universal2.whl", hash = "sha256:9ebafa0b96c62881d5cafa02d9da2e44c23f9f0cd829f3a32a6aff771449c996"}, + {file = "uvloop-0.20.0-cp310-cp310-macosx_10_9_x86_64.whl", hash = "sha256:35968fc697b0527a06e134999eef859b4034b37aebca537daeb598b9d45a137b"}, + {file = "uvloop-0.20.0-cp310-cp310-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:b16696f10e59d7580979b420eedf6650010a4a9c3bd8113f24a103dfdb770b10"}, + {file = "uvloop-0.20.0-cp310-cp310-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:9b04d96188d365151d1af41fa2d23257b674e7ead68cfd61c725a422764062ae"}, + {file = "uvloop-0.20.0-cp310-cp310-musllinux_1_1_aarch64.whl", hash = "sha256:94707205efbe809dfa3a0d09c08bef1352f5d3d6612a506f10a319933757c006"}, + {file = "uvloop-0.20.0-cp310-cp310-musllinux_1_1_x86_64.whl", hash = "sha256:89e8d33bb88d7263f74dc57d69f0063e06b5a5ce50bb9a6b32f5fcbe655f9e73"}, + {file = "uvloop-0.20.0-cp311-cp311-macosx_10_9_universal2.whl", hash = "sha256:e50289c101495e0d1bb0bfcb4a60adde56e32f4449a67216a1ab2750aa84f037"}, + {file = "uvloop-0.20.0-cp311-cp311-macosx_10_9_x86_64.whl", hash = "sha256:e237f9c1e8a00e7d9ddaa288e535dc337a39bcbf679f290aee9d26df9e72bce9"}, + {file = "uvloop-0.20.0-cp311-cp311-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:746242cd703dc2b37f9d8b9f173749c15e9a918ddb021575a0205ec29a38d31e"}, + {file = "uvloop-0.20.0-cp311-cp311-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:82edbfd3df39fb3d108fc079ebc461330f7c2e33dbd002d146bf7c445ba6e756"}, + {file = "uvloop-0.20.0-cp311-cp311-musllinux_1_1_aarch64.whl", hash = "sha256:80dc1b139516be2077b3e57ce1cb65bfed09149e1d175e0478e7a987863b68f0"}, + {file = "uvloop-0.20.0-cp311-cp311-musllinux_1_1_x86_64.whl", hash = "sha256:4f44af67bf39af25db4c1ac27e82e9665717f9c26af2369c404be865c8818dcf"}, + {file = "uvloop-0.20.0-cp312-cp312-macosx_10_9_universal2.whl", hash = "sha256:4b75f2950ddb6feed85336412b9a0c310a2edbcf4cf931aa5cfe29034829676d"}, + {file = "uvloop-0.20.0-cp312-cp312-macosx_10_9_x86_64.whl", hash = "sha256:77fbc69c287596880ecec2d4c7a62346bef08b6209749bf6ce8c22bbaca0239e"}, + {file = "uvloop-0.20.0-cp312-cp312-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:6462c95f48e2d8d4c993a2950cd3d31ab061864d1c226bbf0ee2f1a8f36674b9"}, + {file = "uvloop-0.20.0-cp312-cp312-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:649c33034979273fa71aa25d0fe120ad1777c551d8c4cd2c0c9851d88fcb13ab"}, + {file = "uvloop-0.20.0-cp312-cp312-musllinux_1_1_aarch64.whl", hash = "sha256:3a609780e942d43a275a617c0839d85f95c334bad29c4c0918252085113285b5"}, + {file = "uvloop-0.20.0-cp312-cp312-musllinux_1_1_x86_64.whl", hash = "sha256:aea15c78e0d9ad6555ed201344ae36db5c63d428818b4b2a42842b3870127c00"}, + {file = "uvloop-0.20.0-cp38-cp38-macosx_10_9_universal2.whl", hash = "sha256:f0e94b221295b5e69de57a1bd4aeb0b3a29f61be6e1b478bb8a69a73377db7ba"}, + {file = "uvloop-0.20.0-cp38-cp38-macosx_10_9_x86_64.whl", hash = "sha256:fee6044b64c965c425b65a4e17719953b96e065c5b7e09b599ff332bb2744bdf"}, + {file = "uvloop-0.20.0-cp38-cp38-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:265a99a2ff41a0fd56c19c3838b29bf54d1d177964c300dad388b27e84fd7847"}, + {file = "uvloop-0.20.0-cp38-cp38-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:b10c2956efcecb981bf9cfb8184d27d5d64b9033f917115a960b83f11bfa0d6b"}, + {file = "uvloop-0.20.0-cp38-cp38-musllinux_1_1_aarch64.whl", hash = "sha256:e7d61fe8e8d9335fac1bf8d5d82820b4808dd7a43020c149b63a1ada953d48a6"}, + {file = "uvloop-0.20.0-cp38-cp38-musllinux_1_1_x86_64.whl", hash = "sha256:2beee18efd33fa6fdb0976e18475a4042cd31c7433c866e8a09ab604c7c22ff2"}, + {file = "uvloop-0.20.0-cp39-cp39-macosx_10_9_universal2.whl", hash = "sha256:d8c36fdf3e02cec92aed2d44f63565ad1522a499c654f07935c8f9d04db69e95"}, + {file = "uvloop-0.20.0-cp39-cp39-macosx_10_9_x86_64.whl", hash = "sha256:a0fac7be202596c7126146660725157d4813aa29a4cc990fe51346f75ff8fde7"}, + {file = "uvloop-0.20.0-cp39-cp39-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:9d0fba61846f294bce41eb44d60d58136090ea2b5b99efd21cbdf4e21927c56a"}, + {file = "uvloop-0.20.0-cp39-cp39-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:95720bae002ac357202e0d866128eb1ac82545bcf0b549b9abe91b5178d9b541"}, + {file = "uvloop-0.20.0-cp39-cp39-musllinux_1_1_aarch64.whl", hash = "sha256:36c530d8fa03bfa7085af54a48f2ca16ab74df3ec7108a46ba82fd8b411a2315"}, + {file = "uvloop-0.20.0-cp39-cp39-musllinux_1_1_x86_64.whl", hash = "sha256:e97152983442b499d7a71e44f29baa75b3b02e65d9c44ba53b10338e98dedb66"}, + {file = "uvloop-0.20.0.tar.gz", hash = "sha256:4603ca714a754fc8d9b197e325db25b2ea045385e8a3ad05d3463de725fdf469"}, +] + +[package.extras] +docs = ["Sphinx (>=4.1.2,<4.2.0)", "sphinx-rtd-theme (>=0.5.2,<0.6.0)", "sphinxcontrib-asyncio (>=0.3.0,<0.4.0)"] +test = ["Cython (>=0.29.36,<0.30.0)", "aiohttp (==3.9.0b0)", "aiohttp (>=3.8.1)", "flake8 (>=5.0,<6.0)", "mypy (>=0.800)", "psutil", "pyOpenSSL (>=23.0.0,<23.1.0)", "pycodestyle (>=2.9.0,<2.10.0)"] + [[package]] name = "watchdog" version = "4.0.1" @@ -3082,4 +3126,4 @@ test = ["big-O", "importlib-resources", "jaraco.functools", "jaraco.itertools", [metadata] lock-version = "2.0" python-versions = ">=3.9.0,<4.0" -content-hash = "3c3d4b9ce6b0609d0399ca9aece50495d0a29b978042908dd981ea267f541934" +content-hash = "f4ed3ad03f1ed1db10fcd18edd6665fc691411473026e2ddb5c8634043b6ef96" diff --git a/libs/langgraph/profile.svg b/libs/langgraph/profile.svg new file mode 100644 index 000000000..5832e76ea --- /dev/null +++ b/libs/langgraph/profile.svg @@ -0,0 +1,415 @@ +py-spy record -o profile.svg -- python ../../perf_test.py Reset ZoomSearch submit (concurrent/futures/thread.py:172) (9 samples, 0.14%)submit (concurrent/futures/thread.py:173) (11 samples, 0.17%)_aroute (graph/graph.py:101) (27 samples, 0.42%)to_thread (asyncio/threads.py:25) (27 samples, 0.42%)ainvoke (utils/runnable.py:145) (13 samples, 0.20%)get_async_callback_manager_for_config (langchain_core/runnables/config.py:503) (13 samples, 0.20%)configure (langchain_core/callbacks/manager.py:1996) (13 samples, 0.20%)_configure (langchain_core/callbacks/manager.py:2191) (12 samples, 0.19%)_tracing_v2_is_enabled (langchain_core/tracers/context.py:149) (12 samples, 0.19%)tracing_is_enabled (langsmith/utils.py:105) (12 samples, 0.19%)get_env_var (langsmith/utils.py:376) (10 samples, 0.16%)get (<frozen _collections_abc>:774) (10 samples, 0.16%)__getitem__ (<frozen os>:679) (10 samples, 0.16%)_aroute (graph/graph.py:108) (26 samples, 0.41%)astream (pregel/__init__.py:1397) (9 samples, 0.14%)__aenter__ (pregel/loop.py:698) (9 samples, 0.14%)enter_async_context (contextlib.py:650) (9 samples, 0.14%)__aenter__ (contextlib.py:210) (9 samples, 0.14%)astream (pregel/__init__.py:1422) (10 samples, 0.16%)to_thread (asyncio/threads.py:25) (8 samples, 0.12%)_output_writes (pregel/loop.py:513) (9 samples, 0.14%)_emit (pregel/loop.py:506) (8 samples, 0.12%)<genexpr> (pregel/loop.py:513) (8 samples, 0.12%)atick (pregel/runner.py:131) (16 samples, 0.25%)put_writes (pregel/loop.py:227) (16 samples, 0.25%)_output_writes (pregel/loop.py:518) (7 samples, 0.11%)_emit (pregel/loop.py:506) (7 samples, 0.11%)<genexpr> (pregel/loop.py:518) (7 samples, 0.11%)astream (pregel/__init__.py:1429) (21 samples, 0.33%)ainvoke (pregel/__init__.py:1539) (52 samples, 0.81%)ensure_config (langchain_core/runnables/config.py:174) (14 samples, 0.22%)<dictcomp> (langchain_core/runnables/config.py:175) (14 samples, 0.22%)copy (langchain_core/callbacks/base.py:913) (13 samples, 0.20%)__init__ (langchain_core/callbacks/base.py:907) (7 samples, 0.11%)ainvoke (utils/runnable.py:142) (34 samples, 0.53%)ainvoke (utils/runnable.py:145) (12 samples, 0.19%)get_async_callback_manager_for_config (langchain_core/runnables/config.py:503) (12 samples, 0.19%)configure (langchain_core/callbacks/manager.py:1996) (12 samples, 0.19%)ainvoke (utils/runnable.py:146) (8 samples, 0.12%)ainvoke (utils/runnable.py:154) (7 samples, 0.11%)shield (asyncio/tasks.py:884) (13 samples, 0.20%)_ensure_future (asyncio/tasks.py:680) (13 samples, 0.20%)ainvoke (utils/runnable.py:164) (14 samples, 0.22%)wrapped (langchain_core/callbacks/manager.py:237) (14 samples, 0.22%)ensure_config (langchain_core/runnables/config.py:176) (10 samples, 0.16%)ainvoke (utils/runnable.py:359) (16 samples, 0.25%)tracing_is_enabled (langsmith/utils.py:105) (8 samples, 0.12%)get_env_var (langsmith/utils.py:376) (8 samples, 0.12%)get (<frozen _collections_abc>:774) (8 samples, 0.12%)ainvoke (utils/runnable.py:360) (31 samples, 0.48%)get_async_callback_manager_for_config (langchain_core/runnables/config.py:503) (28 samples, 0.44%)configure (langchain_core/callbacks/manager.py:1996) (28 samples, 0.44%)_configure (langchain_core/callbacks/manager.py:2191) (10 samples, 0.16%)_tracing_v2_is_enabled (langchain_core/tracers/context.py:149) (10 samples, 0.16%)shield (asyncio/tasks.py:884) (10 samples, 0.16%)_ensure_future (asyncio/tasks.py:680) (9 samples, 0.14%)arun_with_retry (pregel/retry.py:79) (81 samples, 1.27%)ainvoke (utils/runnable.py:391) (14 samples, 0.22%)wrapped (langchain_core/callbacks/manager.py:237) (14 samples, 0.22%)<module> (perf_test.py:81) (379 samples, 5.92%)<module>..run (asyncio/runners.py:190) (379 samples, 5.92%)run (asy..run (asyncio/runners.py:118) (379 samples, 5.92%)run (asy.._worker (concurrent/futures/thread.py:81) (5,805 samples, 90.66%)_worker (concurrent/futures/thread.py:81)local_read (pregel/algo.py:109) (9 samples, 0.14%)ChannelsManager (pregel/manager.py:44) (10 samples, 0.16%)__enter__ (contextlib.py:137) (15 samples, 0.23%)local_read (pregel/algo.py:110) (24 samples, 0.37%)__exit__ (contextlib.py:144) (9 samples, 0.14%)ChannelsManager (pregel/manager.py:42) (9 samples, 0.14%)__exit__ (contextlib.py:586) (7 samples, 0.11%)do_read (pregel/read.py:102) (42 samples, 0.66%)local_read (pregel/algo.py:114) (9 samples, 0.14%)do_read (pregel/read.py:86) (8 samples, 0.12%)tick (pregel/loop.py:257) (21 samples, 0.33%)apply_writes (pregel/algo.py:218) (11 samples, 0.17%)update (channels/last_value.py:45) (9 samples, 0.14%)tick (pregel/loop.py:267) (10 samples, 0.16%)_emit (pregel/loop.py:506) (10 samples, 0.16%)<genexpr> (pregel/loop.py:267) (10 samples, 0.16%)map_output_values (pregel/io.py:86) (10 samples, 0.16%)<setcomp> (pregel/io.py:86) (10 samples, 0.16%)tick (pregel/loop.py:285) (8 samples, 0.12%)<genexpr> (pregel/algo.py:383) (16 samples, 0.25%)read_channel (pregel/io.py:19) (16 samples, 0.25%)get (channels/ephemeral_value.py:65) (7 samples, 0.11%)prepare_next_tasks (pregel/algo.py:379) (21 samples, 0.33%)prepare_next_tasks (pregel/algo.py:406) (31 samples, 0.48%)tick (pregel/loop.py:300) (85 samples, 1.33%)prepare_next_tasks (pregel/algo.py:425) (15 samples, 0.23%)run (concurrent/futures/thread.py:58) (206 samples, 3.22%)run..tick (pregel/loop.py:364) (18 samples, 0.28%)_emit (pregel/loop.py:506) (18 samples, 0.28%)<genexpr> (pregel/loop.py:364) (18 samples, 0.28%)map_debug_tasks (pregel/debug.py:94) (13 samples, 0.20%)_worker (concurrent/futures/thread.py:83) (209 samples, 3.26%)_wo..all (6,403 samples, 100%)_bootstrap (threading.py:1002) (6,017 samples, 93.97%)_bootstrap (threading.py:1002)_bootstrap_inner (threading.py:1045) (6,017 samples, 93.97%)_bootstrap_inner (threading.py:1045)run (threading.py:982) (6,017 samples, 93.97%)run (threading.py:982) \ No newline at end of file diff --git a/libs/langgraph/pyproject.toml b/libs/langgraph/pyproject.toml index c32442040..e5e3b7a16 100644 --- a/libs/langgraph/pyproject.toml +++ b/libs/langgraph/pyproject.toml @@ -31,6 +31,7 @@ langgraph-checkpoint = {path = "../checkpoint", develop = true} langgraph-checkpoint-sqlite = {path = "../checkpoint-sqlite", develop = true} langgraph-checkpoint-postgres = {path = "../checkpoint-postgres", develop = true} psycopg = {extras = ["binary"], version = ">=3.0.0"} +uvloop = "^0.20.0" [tool.ruff] lint.select = [ "E", "F", "I" ] diff --git a/perf_test.prof b/perf_test.prof new file mode 100644 index 000000000..7ddbc252c Binary files /dev/null and b/perf_test.prof differ diff --git a/perf_test.py b/perf_test.py new file mode 100644 index 000000000..7c4506525 --- /dev/null +++ b/perf_test.py @@ -0,0 +1,81 @@ +import asyncio +import operator +import random +from time import perf_counter_ns +from typing import Annotated, TypedDict + +import uvloop +from langgraph.checkpoint.memory import MemorySaver +from langgraph.constants import END, START, Send +from langgraph.graph.state import StateGraph + +asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) + + +class OverallState(TypedDict): + subjects: list[str] + jokes: Annotated[list[str], operator.add] + + +async def continue_to_jokes(state: OverallState): + return [Send("generate_joke", {"subject": s}) for s in state["subjects"]] + + +class JokeInput(TypedDict): + subject: str + + +class JokeOutput(TypedDict): + jokes: list[str] + + +async def bump(state: JokeOutput): + return {"jokes": [state["jokes"][0] + " a"]} + + +async def generate(state: JokeInput): + return {"jokes": [f"Joke about {state['subject']}"]} + + +async def edit(state: JokeInput): + subject = state["subject"] + return {"subject": f"{subject} - hohoho"} + + +async def bump_loop(state: JokeOutput): + return END if state["jokes"][0].endswith(" a" * 10) else "bump" + + +# subgraph +subgraph = StateGraph(input=JokeInput, output=JokeOutput) +subgraph.add_node("edit", edit) +subgraph.add_node("generate", generate) +subgraph.add_node("bump", bump) +subgraph.set_entry_point("edit") +subgraph.add_edge("edit", "generate") +subgraph.add_edge("generate", "bump") +subgraph.add_conditional_edges("bump", bump_loop) +subgraph.set_finish_point("generate") +subgraphc = subgraph.compile() + +# parent graph +builder = StateGraph(OverallState) +builder.add_node("generate_joke", subgraphc) +builder.add_conditional_edges(START, continue_to_jokes) +builder.add_edge("generate_joke", END) + + +async def main(): + graph = builder.compile(checkpointer=None) + config = {"configurable": {"thread_id": "1"}} + input = { + "subjects": [random.choice("abcdefghijklmnopqrstuvwxyz") for _ in range(1000)] + } + + # invoke and pause at nested interrupt + s = perf_counter_ns() + len([c async for c in graph.astream(input, config=config)]) + print("Time taken:", (perf_counter_ns() - s) / 1e9) + + +asyncio.run(main()) diff --git a/perf_test_sync.py b/perf_test_sync.py new file mode 100644 index 000000000..ac3e1a07f --- /dev/null +++ b/perf_test_sync.py @@ -0,0 +1,80 @@ +import operator +import random +from time import perf_counter_ns +from typing import Annotated, TypedDict + +from langgraph.checkpoint.sqlite import SqliteSaver +from langgraph.constants import END, START, Send +from langgraph.graph.state import StateGraph + + +class OverallState(TypedDict): + subjects: list[str] + jokes: Annotated[list[str], operator.add] + + +def continue_to_jokes(state: OverallState): + return [Send("generate_joke", {"subject": s}) for s in state["subjects"]] + + +class JokeInput(TypedDict): + subject: str + + +class JokeOutput(TypedDict): + jokes: list[str] + + +def bump(state: JokeOutput): + return {"jokes": [state["jokes"][0] + " a"]} + + +def generate(state: JokeInput): + return {"jokes": [f"Joke about {state['subject']}"]} + + +def edit(state: JokeInput): + subject = state["subject"] + return {"subject": f"{subject} - hohoho"} + + +def bump_loop(state: JokeOutput): + return END if state["jokes"][0].endswith(" a" * 10) else "bump" + + +# subgraph +subgraph = StateGraph(input=JokeInput, output=JokeOutput) +subgraph.add_node("edit", edit) +subgraph.add_node("generate", generate) +subgraph.add_node("bump", bump) +subgraph.set_entry_point("edit") +subgraph.add_edge("edit", "generate") +subgraph.add_edge("generate", "bump") +subgraph.add_conditional_edges("bump", bump_loop) +subgraph.set_finish_point("generate") +subgraphc = subgraph.compile() + +# parent graph +builder = StateGraph(OverallState) +builder.add_node("generate_joke", subgraphc) +builder.add_conditional_edges(START, continue_to_jokes) +builder.add_edge("generate_joke", END) + + +def main(): + with SqliteSaver.from_conn_string(":memory:") as checkpointer: + graph = builder.compile(checkpointer=checkpointer) + config = {"configurable": {"thread_id": "1"}} + input = { + "subjects": [ + random.choice("abcdefghijklmnopqrstuvwxyz") for _ in range(100) + ] + } + + # invoke and pause at nested interrupt + s = perf_counter_ns() + assert len([c for c in graph.stream(input, config=config)]) == 100 + print("Time taken:", (perf_counter_ns() - s) / 1e9) + + +main()