diff --git a/libs/langgraph/langgraph/pregel/__init__.py b/libs/langgraph/langgraph/pregel/__init__.py index 6f93335eb..04e302232 100644 --- a/libs/langgraph/langgraph/pregel/__init__.py +++ b/libs/langgraph/langgraph/pregel/__init__.py @@ -437,8 +437,8 @@ class Pregel(Runnable[Union[dict[str, Any], Any], Union[dict[str, Any], Any]]): saved.checkpoint, LoopProtocol( config=saved.config, - step=saved.metadata["step"], - stop=saved.metadata["step"] + 1, + step=saved.metadata.get("step", -1) + 1, + stop=saved.metadata.get("step", -1) + 2, ), skip_context=True, ) as (channels, managed): @@ -522,8 +522,8 @@ class Pregel(Runnable[Union[dict[str, Any], Any], Union[dict[str, Any], Any]]): saved.checkpoint, LoopProtocol( config=saved.config, - step=saved.metadata["step"], - stop=saved.metadata["step"] + 1, + step=saved.metadata.get("step", -1) + 1, + stop=saved.metadata.get("step", -1) + 2, ), skip_context=True, ) as ( @@ -852,7 +852,7 @@ class Pregel(Runnable[Union[dict[str, Any], Any], Union[dict[str, Any], Any]]): with ChannelsManager( self.channels, checkpoint, - LoopProtocol(config=config, step=step, stop=step + 1), + LoopProtocol(config=config, step=step + 1, stop=step + 2), ) as ( channels, managed, @@ -1002,7 +1002,7 @@ class Pregel(Runnable[Union[dict[str, Any], Any], Union[dict[str, Any], Any]]): async with AsyncChannelsManager( self.channels, checkpoint, - LoopProtocol(config=config, step=step, stop=step + 1), + LoopProtocol(config=config, step=step + 1, stop=step + 2), ) as ( channels, managed, diff --git a/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py b/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py index da2c2bafb..a7c99900d 100644 --- a/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py +++ b/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py @@ -191,7 +191,6 @@ class AsyncKafkaExecutor(AbstractAsyncContextManager): step=saved.metadata["step"] + 1, stop=saved.metadata["step"] + 2, ), - self.graph.store, ) as (channels, managed), AsyncBackgroundExecutor() as submit: if task := await asyncio.to_thread( prepare_single_task,