diff --git a/libs/scheduler-kafka/tests/test_subgraph.py b/libs/scheduler-kafka/tests/test_subgraph.py index 89a092b83..2a6c9992a 100644 --- a/libs/scheduler-kafka/tests/test_subgraph.py +++ b/libs/scheduler-kafka/tests/test_subgraph.py @@ -15,7 +15,7 @@ from langgraph.graph.state import StateGraph from langgraph.pregel import Pregel from langgraph.scheduler.kafka import serde from langgraph.scheduler.kafka.types import MessageToOrchestrator, Topics -from tests.any import AnyDict +from tests.any import AnyDict, AnyInt from tests.drain import drain_topics_async from tests.messages import _AnyIdAIMessage, _AnyIdHumanMessage @@ -198,7 +198,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[0].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": None, "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -265,7 +271,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[0].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -362,7 +374,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[0].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -469,7 +487,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": None, "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"] @@ -531,7 +555,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"] @@ -649,7 +679,13 @@ async def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"] diff --git a/libs/scheduler-kafka/tests/test_subgraph_sync.py b/libs/scheduler-kafka/tests/test_subgraph_sync.py index e7e7bdfb0..c2c9a8fc1 100644 --- a/libs/scheduler-kafka/tests/test_subgraph_sync.py +++ b/libs/scheduler-kafka/tests/test_subgraph_sync.py @@ -15,7 +15,7 @@ from langgraph.pregel import Pregel from langgraph.scheduler.kafka import serde from langgraph.scheduler.kafka.default_sync import DefaultProducer from langgraph.scheduler.kafka.types import MessageToOrchestrator, Topics -from tests.any import AnyDict +from tests.any import AnyDict, AnyInt from tests.drain import drain_topics from tests.messages import _AnyIdAIMessage, _AnyIdHumanMessage @@ -197,7 +197,13 @@ def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[0].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": None, "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -264,7 +270,13 @@ def test_subgraph_w_interrupt( "__pregel_resuming": False, "__pregel_previous": None, "__pregel_task_id": history[0].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -361,7 +373,13 @@ def test_subgraph_w_interrupt( "__pregel_resuming": False, "__pregel_task_id": history[0].tasks[0].id, "__pregel_previous": None, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[0].config["configurable"]["checkpoint_id"] @@ -467,7 +485,13 @@ def test_subgraph_w_interrupt( "__pregel_resuming": True, "__pregel_previous": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": None, "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"] @@ -529,7 +553,13 @@ def test_subgraph_w_interrupt( "__pregel_resuming": True, "__pregel_previous": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"] @@ -647,7 +677,13 @@ def test_subgraph_w_interrupt( "__pregel_previous": None, "__pregel_store": None, "__pregel_task_id": history[1].tasks[0].id, - "__pregel_scratchpad": None, + "__pregel_scratchpad": { + "subgraph_counter": AnyInt(), + "call_counter": 0, + "interrupt_counter": -1, + "null_resume": None, + "resume": [], + }, "checkpoint_id": c.config["configurable"]["checkpoint_id"], "checkpoint_map": { "": history[1].config["configurable"]["checkpoint_id"]