@@ -45,6 +45,7 @@ async def run(self, input):
4545 # for explicit reuse and when an HTTP consumer cancels/closes its run.
4646 streams = {}
4747 self ._tool_streams = streams
48+ self ._response_message_ids = {}
4849 self ._upstream_stream = None
4950 self ._graph_stream = None
5051 try :
@@ -66,6 +67,7 @@ async def run(self, input):
6667 await self ._graph_stream .aclose ()
6768 finally :
6869 streams .clear ()
70+ self ._response_message_ids .clear ()
6971
7072 def _handle_stream_events (self , input ):
7173 # Upstream run uses a bare async-for. Keep its generator so closing
@@ -82,7 +84,25 @@ async def _handle_single_event(self, event, state):
8284 active_run = self .active_run
8385 kind = event .get ("event" )
8486 model_key = (self ._current_lane (), event .get ("run_id" ))
87+ if kind == "on_chat_model_stream" :
88+ chunk = event .get ("data" , {}).get ("chunk" )
89+ # Responses API announces its durable ID in a metadata-only chunk.
90+ # Later chunks get lc_run IDs from LangChain; emitting those creates
91+ # a second copy when the final snapshot restores the provider ID.
92+ metadata = _get (chunk , "response_metadata" , {}) or {}
93+ response_id = metadata .get ("id" )
94+ if response_id and _get (chunk , "id" ) == response_id :
95+ self ._response_message_ids .setdefault (model_key , response_id )
96+ stable_id = self ._response_message_ids .get (model_key )
97+ if stable_id and _get (chunk , "id" ) != stable_id :
98+ clean = (
99+ {** chunk , "id" : stable_id }
100+ if isinstance (chunk , dict )
101+ else chunk .model_copy (update = {"id" : stable_id })
102+ )
103+ event = {** event , "data" : {** event ["data" ], "chunk" : clean }}
85104 if kind == "on_chat_model_end" :
105+ self ._response_message_ids .pop (model_key , None )
86106 saved = self .get_message_in_progress (self .active_run ["id" ])
87107 try :
88108 for slot in self ._tool_streams .pop (model_key , {}).values ():
0 commit comments