From 3bdb7d09beae7e1fcfebe86aa91b18e5a93d2cc1 Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Wed, 14 May 2025 11:30:49 -0700 Subject: [PATCH] Lint --- .../langgraph/scheduler/kafka/executor.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py b/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py index 6564b28a0..9d0af1c8b 100644 --- a/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py +++ b/libs/scheduler-kafka/langgraph/scheduler/kafka/executor.py @@ -221,9 +221,10 @@ class AsyncKafkaExecutor(AbstractAsyncContextManager): runner = PregelRunner( submit=weakref.ref(submit), put_writes=weakref.ref(put_writes), - schedule_task=weakref.WeakMethod(self._schedule_task), ) - async for _ in runner.atick([task], reraise=False): + async for _ in runner.atick( + [task], reraise=False, schedule_task=self._schedule_task + ): pass else: # task was not found @@ -438,9 +439,10 @@ class KafkaExecutor(AbstractContextManager): runner = PregelRunner( submit=weakref.ref(submit), put_writes=weakref.ref(put_writes), - schedule_task=weakref.WeakMethod(self._schedule_task), ) - for _ in runner.tick([task], reraise=False): + for _ in runner.tick( + [task], reraise=False, schedule_task=self._schedule_task + ): pass else: # task was not found