From 58aaffb6bf2af11727d033feab564b33db410381 Mon Sep 17 00:00:00 2001 From: Nuno Campos Date: Mon, 23 Sep 2024 17:34:04 -0700 Subject: [PATCH] Fix edge cases with stream_mode=messages - running callback handler in background thread could potentially lead to ordering issues - deduping on `id()` could lead to messages being dropped if they reused memory address of a previous chunk --- libs/langgraph/langgraph/pregel/messages.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/libs/langgraph/langgraph/pregel/messages.py b/libs/langgraph/langgraph/pregel/messages.py index d0ae539e2..08327805f 100644 --- a/libs/langgraph/langgraph/pregel/messages.py +++ b/libs/langgraph/langgraph/pregel/messages.py @@ -27,21 +27,20 @@ class StreamMessagesHandler(BaseCallbackHandler, _StreamingCallbackHandler): """A callback handler that implements stream_mode=messages. Collects messages from (1) chat model stream events and (2) node outputs.""" + run_inline = True + """We want this callback to run in the main thread, to avoid order/locking issues.""" + def __init__(self, stream: Callable[[StreamChunk], None]): self.stream = stream self.metadata: dict[UUID, Meta] = {} self.seen: set[Union[int, str]] = set() def _emit(self, meta: Meta, message: BaseMessage, *, dedupe: bool = False) -> None: - ident = id(message) if dedupe and message.id in self.seen: return - elif ident in self.seen: - return else: if message.id is None: message.id = str(uuid4()) - self.seen.add(ident) self.seen.add(message.id) self.stream((meta[0], "messages", (message, meta[1])))