From 8207d3fefb30be8736cd902946a81a0828dbabec Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Fri, 23 May 2025 16:11:12 -0700 Subject: [PATCH 1/4] Add tests for stream_events when using imperative api --- libs/langgraph/langgraph/pregel/runner.py | 2 +- libs/langgraph/tests/test_pregel.py | 6 + libs/langgraph/tests/test_pregel_async.py | 270 ++++++++++++++++++++++ libs/langgraph/uv.lock | 5 +- 4 files changed, 280 insertions(+), 3 deletions(-) diff --git a/libs/langgraph/langgraph/pregel/runner.py b/libs/langgraph/langgraph/pregel/runner.py index bbdb90bae..59ede35f5 100644 --- a/libs/langgraph/langgraph/pregel/runner.py +++ b/libs/langgraph/langgraph/pregel/runner.py @@ -457,7 +457,7 @@ def _should_stop_others( if fut.cancelled(): continue elif exc := fut.exception(): - if not isinstance(exc, GraphBubbleUp): + if not isinstance(exc, GraphBubbleUp) and fut not in SKIP_RERAISE_SET: return True return False diff --git a/libs/langgraph/tests/test_pregel.py b/libs/langgraph/tests/test_pregel.py index 81b5eadf8..ac56a4aa2 100644 --- a/libs/langgraph/tests/test_pregel.py +++ b/libs/langgraph/tests/test_pregel.py @@ -8796,3 +8796,9 @@ def test_imp_exception( thread1 = {"configurable": {"thread_id": "1"}} assert my_workflow.invoke(1, thread1) == "done" + + assert [c for c in my_workflow.stream(1, thread1)] == [ + {"my_task": 2}, + {"my_task": 2}, + {"my_workflow": "done"}, + ] diff --git a/libs/langgraph/tests/test_pregel_async.py b/libs/langgraph/tests/test_pregel_async.py index f77612241..739de4312 100644 --- a/libs/langgraph/tests/test_pregel_async.py +++ b/libs/langgraph/tests/test_pregel_async.py @@ -9176,3 +9176,273 @@ async def test_imp_exception( thread1 = {"configurable": {"thread_id": "1"}} assert await my_workflow.ainvoke(1, thread1) == "done" + + assert [c async for c in my_workflow.astream(1, thread1)] == [ + {"my_task": 2}, + {"my_task": 2}, + {"my_workflow": "done"}, + ] + + assert [c async for c in my_workflow.astream_events(1, thread1)] == [ + { + "event": "on_chain_start", + "data": {"input": 1}, + "name": "LangGraph", + "tags": [], + "run_id": AnyStr(), + "metadata": {"thread_id": "1"}, + "parent_ids": [], + }, + { + "event": "on_chain_start", + "data": {"input": 1}, + "name": "my_workflow", + "tags": ["graph:step:4"], + "run_id": AnyStr(), + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_workflow", + "langgraph_triggers": ("__start__",), + "langgraph_path": ("__pregel_pull", "my_workflow"), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [AnyStr()], + }, + { + "event": "on_chain_start", + "data": {"input": {"number": 1}}, + "name": "my_task", + "tags": ["seq:step:1"], + "run_id": AnyStr(), + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "my_task", + "tags": ["seq:step:1"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "data": {"chunk": 2}, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_end", + "data": {"output": 2, "input": {"number": 1}}, + "run_id": AnyStr(), + "name": "my_task", + "tags": ["seq:step:1"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "LangGraph", + "tags": [], + "metadata": {"thread_id": "1"}, + "data": {"chunk": {"my_task": 2}}, + "parent_ids": [], + }, + { + "event": "on_chain_start", + "data": {"input": {"number": 1}}, + "name": "task_with_exception", + "tags": ["seq:step:1"], + "run_id": AnyStr(), + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_start", + "data": {"input": {"number": 1}}, + "name": "my_task", + "tags": ["seq:step:1"], + "run_id": AnyStr(), + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "my_task", + "tags": ["seq:step:1"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "data": {"chunk": 2}, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_end", + "data": {"output": 2, "input": {"number": 1}}, + "run_id": AnyStr(), + "name": "my_task", + "tags": ["seq:step:1"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_task", + "langgraph_triggers": ("__pregel_push",), + "langgraph_path": ( + "__pregel_push", + ("__pregel_pull", "my_workflow"), + 2, + True, + ), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [ + AnyStr(), + AnyStr(), + ], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "my_workflow", + "tags": ["graph:step:4"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_workflow", + "langgraph_triggers": ("__start__",), + "langgraph_path": ("__pregel_pull", "my_workflow"), + "langgraph_checkpoint_ns": AnyStr(), + }, + "data": {"chunk": "done"}, + "parent_ids": [AnyStr()], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "LangGraph", + "tags": [], + "metadata": {"thread_id": "1"}, + "data": {"chunk": {"my_task": 2}}, + "parent_ids": [], + }, + { + "event": "on_chain_end", + "data": {"output": "done", "input": 1}, + "run_id": AnyStr(), + "name": "my_workflow", + "tags": ["graph:step:4"], + "metadata": { + "thread_id": "1", + "langgraph_step": 4, + "langgraph_node": "my_workflow", + "langgraph_triggers": ("__start__",), + "langgraph_path": ("__pregel_pull", "my_workflow"), + "langgraph_checkpoint_ns": AnyStr(), + }, + "parent_ids": [AnyStr()], + }, + { + "event": "on_chain_stream", + "run_id": AnyStr(), + "name": "LangGraph", + "tags": [], + "metadata": {"thread_id": "1"}, + "data": {"chunk": {"my_workflow": "done"}}, + "parent_ids": [], + }, + { + "event": "on_chain_end", + "data": {"output": "done"}, + "run_id": AnyStr(), + "name": "LangGraph", + "tags": [], + "metadata": {"thread_id": "1"}, + "parent_ids": [], + }, + ] diff --git a/libs/langgraph/uv.lock b/libs/langgraph/uv.lock index 793c7a168..0b3ff883c 100644 --- a/libs/langgraph/uv.lock +++ b/libs/langgraph/uv.lock @@ -1,5 +1,4 @@ version = 1 -revision = 1 requires-python = ">=3.9" resolution-markers = [ "python_full_version >= '3.13' and python_full_version < '4.0'", @@ -1197,7 +1196,7 @@ wheels = [ [[package]] name = "langgraph" -version = "0.4.5" +version = "0.4.6" source = { editable = "." } dependencies = [ { name = "langchain-core" }, @@ -2123,6 +2122,8 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/b8/2a/25e0be2b509c28375c7f75c7e8d8d060773f2cce4856a1654276e3202339/pycryptodome-3.22.0-cp37-abi3-musllinux_1_2_x86_64.whl", hash = "sha256:d21c1eda2f42211f18a25db4eaf8056c94a8563cd39da3683f89fe0d881fb772", size = 2262255 }, { url = "https://files.pythonhosted.org/packages/41/58/60917bc4bbd91712e53ce04daf237a74a0ad731383a01288130672994328/pycryptodome-3.22.0-cp37-abi3-win32.whl", hash = "sha256:f02baa9f5e35934c6e8dcec91fcde96612bdefef6e442813b8ea34e82c84bbfb", size = 1763403 }, { url = "https://files.pythonhosted.org/packages/55/f4/244c621afcf7867e23f63cfd7a9630f14cfe946c9be7e566af6c3915bcde/pycryptodome-3.22.0-cp37-abi3-win_amd64.whl", hash = "sha256:d086aed307e96d40c23c42418cbbca22ecc0ab4a8a0e24f87932eeab26c08627", size = 1794568 }, + { url = "https://files.pythonhosted.org/packages/cd/13/16d3a83b07f949a686f6cfd7cfc60e57a769ff502151ea140ad67b118e26/pycryptodome-3.22.0-pp27-pypy_73-manylinux2010_x86_64.whl", hash = "sha256:98fd9da809d5675f3a65dcd9ed384b9dc67edab6a4cda150c5870a8122ec961d", size = 1700779 }, + { url = "https://files.pythonhosted.org/packages/13/af/16d26f7dfc5fd7696ea2c91448f937b51b55312b5bed44f777563e32a4fe/pycryptodome-3.22.0-pp27-pypy_73-win32.whl", hash = "sha256:37ddcd18284e6b36b0a71ea495a4c4dca35bb09ccc9bfd5b91bfaf2321f131c1", size = 1775230 }, { url = "https://files.pythonhosted.org/packages/37/c3/e3423e72669ca09f141aae493e1feaa8b8475859898b04f57078280a61c4/pycryptodome-3.22.0-pp310-pypy310_pp73-macosx_10_15_x86_64.whl", hash = "sha256:b4bdce34af16c1dcc7f8c66185684be15f5818afd2a82b75a4ce6b55f9783e13", size = 1618698 }, { url = "https://files.pythonhosted.org/packages/f9/b7/35eec0b3919cafea362dcb68bb0654d9cb3cde6da6b7a9d8480ce0bf203a/pycryptodome-3.22.0-pp310-pypy310_pp73-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:2988ffcd5137dc2d27eb51cd18c0f0f68e5b009d5fec56fbccb638f90934f333", size = 1666957 }, { url = "https://files.pythonhosted.org/packages/b0/1f/f49bccdd8d61f1da4278eb0d6aee7f988f1a6ec4056b0c2dc51eda45ae27/pycryptodome-3.22.0-pp310-pypy310_pp73-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:e653519dedcd1532788547f00eeb6108cc7ce9efdf5cc9996abce0d53f95d5a9", size = 1659242 }, From 9170f636d0a70f41b45b1061e893fd618f6549a1 Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Fri, 23 May 2025 16:14:15 -0700 Subject: [PATCH 2/4] Lock --- libs/prebuilt/uv.lock | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/libs/prebuilt/uv.lock b/libs/prebuilt/uv.lock index 8ff05ba40..b6982f70e 100644 --- a/libs/prebuilt/uv.lock +++ b/libs/prebuilt/uv.lock @@ -1,5 +1,4 @@ version = 1 -revision = 1 requires-python = ">=3.9" resolution-markers = [ "python_full_version >= '3.12.4'", @@ -320,7 +319,7 @@ wheels = [ [[package]] name = "langgraph" -version = "0.4.5" +version = "0.4.6" source = { editable = "../langgraph" } dependencies = [ { name = "langchain-core" }, From 126a8f5bc612c90df5960748b534a29cab6b38ea Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Fri, 23 May 2025 16:22:12 -0700 Subject: [PATCH 3/4] Lock --- libs/scheduler-kafka/uv.lock | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/libs/scheduler-kafka/uv.lock b/libs/scheduler-kafka/uv.lock index 24f804d6c..0af83c5d0 100644 --- a/libs/scheduler-kafka/uv.lock +++ b/libs/scheduler-kafka/uv.lock @@ -1,5 +1,4 @@ version = 1 -revision = 1 requires-python = ">=3.9" resolution-markers = [ "python_full_version >= '3.12.4'", @@ -447,7 +446,7 @@ wheels = [ [[package]] name = "langgraph" -version = "0.4.5" +version = "0.4.6" source = { editable = "../langgraph" } dependencies = [ { name = "langchain-core" }, From 9bf672835453a52a72f8b7a00be8d08d79923b6d Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Fri, 23 May 2025 16:32:51 -0700 Subject: [PATCH 4/4] Try to make test less flaky --- libs/langgraph/tests/test_pregel.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libs/langgraph/tests/test_pregel.py b/libs/langgraph/tests/test_pregel.py index ac56a4aa2..d5c02a911 100644 --- a/libs/langgraph/tests/test_pregel.py +++ b/libs/langgraph/tests/test_pregel.py @@ -6873,7 +6873,7 @@ def test_sync_streaming_with_functional_api() -> None: should be greater than the time delay between the two tasks. """ - time_delay = 0.01 + time_delay = 0.05 @task() def slow() -> dict: