From 74f3ed511f4cef7d4ef905c65341179a9b10c165 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Tue, 8 Sep 2026 21:10:06 +0800 Subject: [PATCH 1/8] fix(ag-ui): dedupe client-replayed transcripts on resume (#8140) Route non-service-session (and workflow) resume seeding through _reconstruct_messages_from_thread_snapshot so a reference AG-UI client that replays its transcript is not double-persisted into the thread snapshot. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 8 ++- .../ag-ui/agent_framework_ag_ui/_workflow.py | 15 ++++-- .../ag_ui/test_resume_transcript_dedupe.py | 52 +++++++++++++++++++ 3 files changed, 70 insertions(+), 5 deletions(-) create mode 100644 python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index 4fdaef188d5..74858e6c710 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2483,7 +2483,13 @@ async def run_agent_stream( seeded_resume_from_snapshot = True if not config.use_service_session: - raw_messages = snapshot_session.resume_seeded_messages(raw_messages) + # Use the same overlap/reconstructor as non-resume turns so a client + # that replays its transcript on resume is not double-persisted (#8140). + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py index 73b2a1d7ec0..42a1574e168 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py @@ -446,10 +446,17 @@ async def run(self, input_data: dict[str, Any]) -> AsyncGenerator[BaseEvent]: run_checkpoint_storage = _OwnedWorkflowCheckpointStorage(checkpoint_storage, request_owner) builder_seed_messages = raw_messages if resume_payload is not None or (checkpoint_id is not None and not raw_messages): - # Resume requests carry only the synthesized interrupt response, and a - # checkpoint-only resume carries no new messages at all; in both cases seed - # the builder with stored history to avoid persisting a truncated thread. - builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) + # Resume / checkpoint-only requests need stored history. Prefer the + # overlap reconstructor when a snapshot exists so client-replayed + # transcripts are not duplicated (#8140); otherwise prepend stored. + if stored_snapshot is not None: + builder_seed_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=builder_seed_messages, + stored_interrupt=stored_snapshot.interrupt, + ) + else: + builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) snapshot_builder = _WorkflowSnapshotBuilder(builder_seed_messages) if snapshot_session.enabled else None if snapshot_builder is not None and effective_state: # Seed builder state so a run that emits no StateSnapshotEvent still diff --git a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py new file mode 100644 index 00000000000..a23ea62258e --- /dev/null +++ b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py @@ -0,0 +1,52 @@ +"""Regression for AG-UI resume with client-replayed transcript (#8140).""" + +from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot + + +def test_resume_with_replayed_transcript_does_not_duplicate_history(): + """Clients that re-send the full transcript on resume must not double-write.""" + stored = [ + {"id": "u1", "role": "user", "content": "please run the tool"}, + { + "id": "a1", + "role": "assistant", + "content": "", + "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], + }, + ] + # Reference AG-UI client replays the same transcript plus the resume turn. + incoming = [ + *stored, + {"id": "u2", "role": "user", "content": "approved"}, + ] + + reconstructed = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored, + incoming_messages=incoming, + stored_interrupt=[{"interruptId": "int-1"}], + ) + + assert [m.get("id") for m in reconstructed] == ["u1", "a1", "u2"] + assert len(reconstructed) == 3 + + +def test_resume_with_interrupt_only_still_seeds_history(): + """Resume payloads that carry only the synthetic response still get history.""" + stored = [ + {"id": "u1", "role": "user", "content": "please run the tool"}, + { + "id": "a1", + "role": "assistant", + "content": "", + "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], + }, + ] + incoming = [{"id": "u2", "role": "user", "content": "approved"}] + + reconstructed = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored, + incoming_messages=incoming, + stored_interrupt=[{"interruptId": "int-1"}], + ) + + assert [m.get("id") for m in reconstructed] == ["u1", "a1", "u2"] From 3bb98c182bdd58ba7e199cc945ab47e52b2b17ae Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Tue, 8 Sep 2026 23:13:08 +0800 Subject: [PATCH 2/8] fix(ag-ui): keep resume seeding for empty messages Use the transcript reconstructor only when the client sends a non-empty messages list; empty approval/checkpoint resumes still go through resume_seeded_messages so stored history is preserved. Also add the required copyright header and document the empty-input reconstructor contract in tests. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 15 ++++++--- .../ag-ui/agent_framework_ag_ui/_workflow.py | 8 +++-- .../ag_ui/test_resume_transcript_dedupe.py | 33 +++++++++++++++++-- 3 files changed, 46 insertions(+), 10 deletions(-) diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index 74858e6c710..85eb5262770 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2485,11 +2485,16 @@ async def run_agent_stream( if not config.use_service_session: # Use the same overlap/reconstructor as non-resume turns so a client # that replays its transcript on resume is not double-persisted (#8140). - raw_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=raw_messages, - stored_interrupt=stored_snapshot.interrupt, - ) + # Empty messages (interrupt-only / approval resume) still need the + # stored history prepended — the reconstructor returns [] unchanged. + if raw_messages: + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) + else: + raw_messages = snapshot_session.resume_seeded_messages(raw_messages) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py index 42a1574e168..1e0414ed057 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py @@ -447,9 +447,11 @@ async def run(self, input_data: dict[str, Any]) -> AsyncGenerator[BaseEvent]: builder_seed_messages = raw_messages if resume_payload is not None or (checkpoint_id is not None and not raw_messages): # Resume / checkpoint-only requests need stored history. Prefer the - # overlap reconstructor when a snapshot exists so client-replayed - # transcripts are not duplicated (#8140); otherwise prepend stored. - if stored_snapshot is not None: + # overlap reconstructor when a snapshot exists and the client sent a + # (possibly replayed) transcript so it is not duplicated (#8140). + # Empty input must keep resume_seeded_messages — the reconstructor + # returns [] unchanged and would otherwise wipe stored history. + if stored_snapshot is not None and builder_seed_messages: builder_seed_messages = _reconstruct_messages_from_thread_snapshot( stored_messages=stored_snapshot.messages, incoming_messages=builder_seed_messages, diff --git a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py index a23ea62258e..53b12695ca2 100644 --- a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py +++ b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py @@ -1,3 +1,5 @@ +# Copyright (c) Microsoft. All rights reserved. + """Regression for AG-UI resume with client-replayed transcript (#8140).""" from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot @@ -30,8 +32,8 @@ def test_resume_with_replayed_transcript_does_not_duplicate_history(): assert len(reconstructed) == 3 -def test_resume_with_interrupt_only_still_seeds_history(): - """Resume payloads that carry only the synthetic response still get history.""" +def test_resume_with_partial_new_turn_still_seeds_history(): + """Non-empty resume suffix is overlapped onto stored history without duplication.""" stored = [ {"id": "u1", "role": "user", "content": "please run the tool"}, { @@ -50,3 +52,30 @@ def test_resume_with_interrupt_only_still_seeds_history(): ) assert [m.get("id") for m in reconstructed] == ["u1", "a1", "u2"] + + +def test_reconstruct_with_empty_incoming_returns_empty(): + """Empty incoming is unchanged; call sites must use resume_seeded_messages instead. + + Approval / checkpoint-only resumes commonly send ``messages: []``. The + reconstructor intentionally returns that empty list so callers can distinguish + \"no new turn\" from \"replayed transcript\" and fall back to seeding. + """ + stored = [ + {"id": "u1", "role": "user", "content": "please run the tool"}, + { + "id": "a1", + "role": "assistant", + "content": "", + "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], + }, + ] + + assert ( + _reconstruct_messages_from_thread_snapshot( + stored_messages=stored, + incoming_messages=[], + stored_interrupt=[{"interruptId": "int-1"}], + ) + == [] + ) From 4b66801a6b82db5900d4cf04e1066ad8fcdcd3b6 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Wed, 9 Sep 2026 11:38:00 +0800 Subject: [PATCH 3/8] fix(ag-ui): centralize resume message reconcile on ThreadSnapshotSession Add reconcile_resume_messages for empty vs replayed resume shapes, use it from agent and workflow runners, and mark generic reconstructed resumes as seeded so save-time prepend does not duplicate history. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 30 ++++----- .../_snapshot_session.py | 19 +++++- .../ag-ui/agent_framework_ag_ui/_workflow.py | 16 +---- .../ag_ui/test_resume_transcript_dedupe.py | 67 +++++++++++-------- 4 files changed, 73 insertions(+), 59 deletions(-) diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index 85eb5262770..79466a844f9 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2483,18 +2483,8 @@ async def run_agent_stream( seeded_resume_from_snapshot = True if not config.use_service_session: - # Use the same overlap/reconstructor as non-resume turns so a client - # that replays its transcript on resume is not double-persisted (#8140). - # Empty messages (interrupt-only / approval resume) still need the - # stored history prepended — the reconstructor returns [] unchanged. - if raw_messages: - raw_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=raw_messages, - stored_interrupt=stored_snapshot.interrupt, - ) - else: - raw_messages = snapshot_session.resume_seeded_messages(raw_messages) + # One session-owned reconcile for empty vs replayed resume shapes (#8140). + raw_messages = snapshot_session.reconcile_resume_messages(raw_messages) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, @@ -2503,11 +2493,17 @@ async def run_agent_stream( ) raw_messages = provider_suffix elif not config.use_service_session: - raw_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=raw_messages, - stored_interrupt=stored_snapshot.interrupt, - ) + if resume_payload is not None: + # Predictive-state / generic resumes also merge history here; mark seeded + # so save-time resume_seeded_messages does not prepend again (#8140). + seeded_resume_from_snapshot = True + raw_messages = snapshot_session.reconcile_resume_messages(raw_messages) + else: + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py index b0694596e74..24c083980fc 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py @@ -22,7 +22,7 @@ StateSnapshotEvent, ) -from ._run_common import _build_run_finished_event +from ._run_common import _build_run_finished_event, _reconstruct_messages_from_thread_snapshot from ._snapshots import ( AGUIThreadSnapshot, AGUIThreadSnapshotStore, @@ -141,6 +141,23 @@ def resume_seeded_messages(self, incoming: list[dict[str, Any]]) -> list[dict[st return incoming return [copy.deepcopy(message) for message in self._stored.messages] + incoming + def reconcile_resume_messages(self, incoming: list[dict[str, Any]]) -> list[dict[str, Any]]: + """Merge stored history with resume messages without double-persisting. + + Empty incoming (interrupt-only / checkpoint resume) prepends stored + history. Non-empty incoming is overlapped via the thread-snapshot + reconstructor so a client that replays its transcript is not duplicated. + """ + if self._stored is None: + return incoming + if not incoming: + return self.resume_seeded_messages(incoming) + return _reconstruct_messages_from_thread_snapshot( + stored_messages=self._stored.messages, + incoming_messages=incoming, + stored_interrupt=self._stored.interrupt, + ) + async def save( self, *, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py index 1e0414ed057..c1f1540a7b2 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py @@ -446,19 +446,9 @@ async def run(self, input_data: dict[str, Any]) -> AsyncGenerator[BaseEvent]: run_checkpoint_storage = _OwnedWorkflowCheckpointStorage(checkpoint_storage, request_owner) builder_seed_messages = raw_messages if resume_payload is not None or (checkpoint_id is not None and not raw_messages): - # Resume / checkpoint-only requests need stored history. Prefer the - # overlap reconstructor when a snapshot exists and the client sent a - # (possibly replayed) transcript so it is not duplicated (#8140). - # Empty input must keep resume_seeded_messages — the reconstructor - # returns [] unchanged and would otherwise wipe stored history. - if stored_snapshot is not None and builder_seed_messages: - builder_seed_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=builder_seed_messages, - stored_interrupt=stored_snapshot.interrupt, - ) - else: - builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) + # Resume / checkpoint-only requests need stored history. Session-owned + # reconcile covers empty vs client-replayed transcripts (#8140). + builder_seed_messages = snapshot_session.reconcile_resume_messages(builder_seed_messages) snapshot_builder = _WorkflowSnapshotBuilder(builder_seed_messages) if snapshot_session.enabled else None if snapshot_builder is not None and effective_state: # Seed builder state so a run that emits no StateSnapshotEvent still diff --git a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py index 53b12695ca2..693353e2bf8 100644 --- a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py +++ b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py @@ -3,11 +3,12 @@ """Regression for AG-UI resume with client-replayed transcript (#8140).""" from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot +from agent_framework_ag_ui._snapshot_session import ThreadSnapshotSession +from agent_framework_ag_ui._snapshots import AGUIThreadSnapshot -def test_resume_with_replayed_transcript_does_not_duplicate_history(): - """Clients that re-send the full transcript on resume must not double-write.""" - stored = [ +def _stored_history() -> list[dict]: + return [ {"id": "u1", "role": "user", "content": "please run the tool"}, { "id": "a1", @@ -16,7 +17,11 @@ def test_resume_with_replayed_transcript_does_not_duplicate_history(): "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], }, ] - # Reference AG-UI client replays the same transcript plus the resume turn. + + +def test_resume_with_replayed_transcript_does_not_duplicate_history(): + """Clients that re-send the full transcript on resume must not double-write.""" + stored = _stored_history() incoming = [ *stored, {"id": "u2", "role": "user", "content": "approved"}, @@ -34,15 +39,7 @@ def test_resume_with_replayed_transcript_does_not_duplicate_history(): def test_resume_with_partial_new_turn_still_seeds_history(): """Non-empty resume suffix is overlapped onto stored history without duplication.""" - stored = [ - {"id": "u1", "role": "user", "content": "please run the tool"}, - { - "id": "a1", - "role": "assistant", - "content": "", - "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], - }, - ] + stored = _stored_history() incoming = [{"id": "u2", "role": "user", "content": "approved"}] reconstructed = _reconstruct_messages_from_thread_snapshot( @@ -55,21 +52,8 @@ def test_resume_with_partial_new_turn_still_seeds_history(): def test_reconstruct_with_empty_incoming_returns_empty(): - """Empty incoming is unchanged; call sites must use resume_seeded_messages instead. - - Approval / checkpoint-only resumes commonly send ``messages: []``. The - reconstructor intentionally returns that empty list so callers can distinguish - \"no new turn\" from \"replayed transcript\" and fall back to seeding. - """ - stored = [ - {"id": "u1", "role": "user", "content": "please run the tool"}, - { - "id": "a1", - "role": "assistant", - "content": "", - "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], - }, - ] + """Empty incoming is unchanged; call sites must use resume seeding instead.""" + stored = _stored_history() assert ( _reconstruct_messages_from_thread_snapshot( @@ -79,3 +63,30 @@ def test_reconstruct_with_empty_incoming_returns_empty(): ) == [] ) + + +def test_reconcile_resume_messages_empty_vs_replayed(): + """Session-owned reconcile covers interrupt-only and replayed transcript shapes.""" + stored = _stored_history() + session = ThreadSnapshotSession( + store=None, + scope=None, + thread_id="t1", + stored=AGUIThreadSnapshot( + messages=stored, + state=None, + interrupt=[{"interruptId": "int-1"}], + session_state=None, + ), + ) + + empty = session.reconcile_resume_messages([]) + assert [m.get("id") for m in empty] == ["u1", "a1"] + + replayed = session.reconcile_resume_messages( + [ + *stored, + {"id": "u2", "role": "user", "content": "approved"}, + ] + ) + assert [m.get("id") for m in replayed] == ["u1", "a1", "u2"] From e6820218139be57f3fefff0e3a533e1456e6d28b Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Wed, 9 Sep 2026 17:28:25 +0800 Subject: [PATCH 4/8] fix(ag-ui): keep empty confirm_changes resumes unseeded MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Generic/predictive resumes with messages:[] must not prepend stored history into the turn input — confirm_changes synthesizes only the resume tool message and short-circuits without another model call. Seeding that history caused approval mismatch warnings and an extra stream invocation (test_agent_endpoint_confirm_changes_clears_persisted_interrupt). Only non-empty client-replayed generic resumes go through reconcile_resume_messages + seeded_resume_from_snapshot. Approval and workflow/checkpoint paths still seed empty via reconcile. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 9 ++++++--- .../ag-ui/agent_framework_ag_ui/_snapshot_session.py | 11 ++++++++--- 2 files changed, 14 insertions(+), 6 deletions(-) diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index 79466a844f9..807fa304881 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2493,9 +2493,12 @@ async def run_agent_stream( ) raw_messages = provider_suffix elif not config.use_service_session: - if resume_payload is not None: - # Predictive-state / generic resumes also merge history here; mark seeded - # so save-time resume_seeded_messages does not prepend again (#8140). + if resume_payload is not None and raw_messages: + # Client-replayed transcript on predictive/generic resume: overlap-merge + # and mark seeded so save-time resume_seeded_messages does not prepend + # again (#8140). Empty interrupt-only resumes (e.g. confirm_changes) + # must stay empty so synthesized resume tool messages are the only + # turn input; history is restored at save when this flag stays false. seeded_resume_from_snapshot = True raw_messages = snapshot_session.reconcile_resume_messages(raw_messages) else: diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py index 24c083980fc..80232cfdd71 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py @@ -144,9 +144,14 @@ def resume_seeded_messages(self, incoming: list[dict[str, Any]]) -> list[dict[st def reconcile_resume_messages(self, incoming: list[dict[str, Any]]) -> list[dict[str, Any]]: """Merge stored history with resume messages without double-persisting. - Empty incoming (interrupt-only / checkpoint resume) prepends stored - history. Non-empty incoming is overlapped via the thread-snapshot - reconstructor so a client that replays its transcript is not duplicated. + Empty incoming prepends stored history (approval / checkpoint builders that + need history in the provider or snapshot input). Non-empty incoming is + overlapped via the thread-snapshot reconstructor so a client that replays + its transcript is not duplicated. + + Agent confirm_changes resumes that synthesize the tool result separately + should keep empty incoming as empty and call this only for non-empty + client-replayed transcripts. """ if self._stored is None: return incoming From 40a0733b2395204ba1deb8a8e1dbec6a852dce14 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Wed, 9 Sep 2026 22:45:01 +0800 Subject: [PATCH 5/8] refactor(ag-ui): fold reconcile_resume_messages into resume_seeded_messages Address review: empty incoming still prepends stored history; non-empty uses the overlap reconstructor so client-replayed transcripts are not duplicated. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 6 +++--- .../agent_framework_ag_ui/_snapshot_session.py | 16 +++------------- .../ag-ui/agent_framework_ag_ui/_workflow.py | 4 ++-- .../tests/ag_ui/test_resume_transcript_dedupe.py | 10 +++++----- 4 files changed, 13 insertions(+), 23 deletions(-) diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index 807fa304881..c9388e9441d 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2483,8 +2483,8 @@ async def run_agent_stream( seeded_resume_from_snapshot = True if not config.use_service_session: - # One session-owned reconcile for empty vs replayed resume shapes (#8140). - raw_messages = snapshot_session.reconcile_resume_messages(raw_messages) + # Session-owned resume seeding covers empty vs replayed shapes (#8140). + raw_messages = snapshot_session.resume_seeded_messages(raw_messages) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, @@ -2500,7 +2500,7 @@ async def run_agent_stream( # must stay empty so synthesized resume tool messages are the only # turn input; history is restored at save when this flag stays false. seeded_resume_from_snapshot = True - raw_messages = snapshot_session.reconcile_resume_messages(raw_messages) + raw_messages = snapshot_session.resume_seeded_messages(raw_messages) else: raw_messages = _reconstruct_messages_from_thread_snapshot( stored_messages=stored_snapshot.messages, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py index 80232cfdd71..4c2e08c6492 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py @@ -132,17 +132,7 @@ def effective_state( return state def resume_seeded_messages(self, incoming: list[dict[str, Any]]) -> list[dict[str, Any]]: - """Prepend copies of stored thread history to a resume request's messages. - - Resume requests carry only the synthesized interrupt response; seeding - with stored history keeps the persisted thread from being truncated. - """ - if self._stored is None: - return incoming - return [copy.deepcopy(message) for message in self._stored.messages] + incoming - - def reconcile_resume_messages(self, incoming: list[dict[str, Any]]) -> list[dict[str, Any]]: - """Merge stored history with resume messages without double-persisting. + """Merge stored thread history with resume messages without double-persisting. Empty incoming prepends stored history (approval / checkpoint builders that need history in the provider or snapshot input). Non-empty incoming is @@ -151,12 +141,12 @@ def reconcile_resume_messages(self, incoming: list[dict[str, Any]]) -> list[dict Agent confirm_changes resumes that synthesize the tool result separately should keep empty incoming as empty and call this only for non-empty - client-replayed transcripts. + client-replayed transcripts (or at save-time for interrupt-only resumes). """ if self._stored is None: return incoming if not incoming: - return self.resume_seeded_messages(incoming) + return [copy.deepcopy(message) for message in self._stored.messages] return _reconstruct_messages_from_thread_snapshot( stored_messages=self._stored.messages, incoming_messages=incoming, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py index c1f1540a7b2..7580ca77852 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py @@ -447,8 +447,8 @@ async def run(self, input_data: dict[str, Any]) -> AsyncGenerator[BaseEvent]: builder_seed_messages = raw_messages if resume_payload is not None or (checkpoint_id is not None and not raw_messages): # Resume / checkpoint-only requests need stored history. Session-owned - # reconcile covers empty vs client-replayed transcripts (#8140). - builder_seed_messages = snapshot_session.reconcile_resume_messages(builder_seed_messages) + # resume_seeded_messages covers empty vs client-replayed transcripts (#8140). + builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) snapshot_builder = _WorkflowSnapshotBuilder(builder_seed_messages) if snapshot_session.enabled else None if snapshot_builder is not None and effective_state: # Seed builder state so a run that emits no StateSnapshotEvent still diff --git a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py index 693353e2bf8..5ac120cf2ba 100644 --- a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py +++ b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py @@ -52,7 +52,7 @@ def test_resume_with_partial_new_turn_still_seeds_history(): def test_reconstruct_with_empty_incoming_returns_empty(): - """Empty incoming is unchanged; call sites must use resume seeding instead.""" + """Empty incoming is unchanged at the helper; ThreadSnapshotSession.resume_seeded_messages seeds instead.""" stored = _stored_history() assert ( @@ -65,8 +65,8 @@ def test_reconstruct_with_empty_incoming_returns_empty(): ) -def test_reconcile_resume_messages_empty_vs_replayed(): - """Session-owned reconcile covers interrupt-only and replayed transcript shapes.""" +def test_resume_seeded_messages_empty_vs_replayed(): + """Session-owned resume seeding covers interrupt-only and replayed transcript shapes.""" stored = _stored_history() session = ThreadSnapshotSession( store=None, @@ -80,10 +80,10 @@ def test_reconcile_resume_messages_empty_vs_replayed(): ), ) - empty = session.reconcile_resume_messages([]) + empty = session.resume_seeded_messages([]) assert [m.get("id") for m in empty] == ["u1", "a1"] - replayed = session.reconcile_resume_messages( + replayed = session.resume_seeded_messages( [ *stored, {"id": "u2", "role": "user", "content": "approved"}, From 4b0f547c1157f08fd62984e9174b2303c43e73d7 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Thu, 10 Sep 2026 07:09:40 +0800 Subject: [PATCH 6/8] fix(ag-ui): restore resume_seeded prepend; fold #8140 tests Per review: keep resume_seeded_messages as blind prepend (save-time / empty seeds). Use reconstruct only for non-empty client-replayed transcripts. Remove dedicated small test file and fold coverage into test_snapshot_session. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 18 +++- .../_snapshot_session.py | 22 ++--- .../ag-ui/agent_framework_ag_ui/_workflow.py | 13 ++- .../ag_ui/test_resume_transcript_dedupe.py | 92 ------------------- .../tests/ag_ui/test_snapshot_session.py | 29 ++++++ 5 files changed, 62 insertions(+), 112 deletions(-) delete mode 100644 python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index c9388e9441d..cc3274ee9ca 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2483,8 +2483,16 @@ async def run_agent_stream( seeded_resume_from_snapshot = True if not config.use_service_session: - # Session-owned resume seeding covers empty vs replayed shapes (#8140). - raw_messages = snapshot_session.resume_seeded_messages(raw_messages) + # Empty approval resumes prepend stored history; non-empty/replayed + # transcripts overlap-merge so clients do not double-write (#8140). + if raw_messages: + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) + else: + raw_messages = snapshot_session.resume_seeded_messages(raw_messages) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages, @@ -2500,7 +2508,11 @@ async def run_agent_stream( # must stay empty so synthesized resume tool messages are the only # turn input; history is restored at save when this flag stays false. seeded_resume_from_snapshot = True - raw_messages = snapshot_session.resume_seeded_messages(raw_messages) + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) else: raw_messages = _reconstruct_messages_from_thread_snapshot( stored_messages=stored_snapshot.messages, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py index 4c2e08c6492..44cffa4cc72 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_snapshot_session.py @@ -22,7 +22,7 @@ StateSnapshotEvent, ) -from ._run_common import _build_run_finished_event, _reconstruct_messages_from_thread_snapshot +from ._run_common import _build_run_finished_event from ._snapshots import ( AGUIThreadSnapshot, AGUIThreadSnapshotStore, @@ -132,26 +132,20 @@ def effective_state( return state def resume_seeded_messages(self, incoming: list[dict[str, Any]]) -> list[dict[str, Any]]: - """Merge stored thread history with resume messages without double-persisting. + """Prepend copies of stored thread history to a resume request's messages. - Empty incoming prepends stored history (approval / checkpoint builders that - need history in the provider or snapshot input). Non-empty incoming is - overlapped via the thread-snapshot reconstructor so a client that replays - its transcript is not duplicated. + Resume requests often carry only the synthesized interrupt response; seeding + with stored history keeps the persisted thread from being truncated. - Agent confirm_changes resumes that synthesize the tool result separately - should keep empty incoming as empty and call this only for non-empty - client-replayed transcripts (or at save-time for interrupt-only resumes). + For non-empty client-replayed transcripts that already overlap stored + history, callers should use ``_reconstruct_messages_from_thread_snapshot`` + instead so messages are not double-persisted (#8140). """ if self._stored is None: return incoming if not incoming: return [copy.deepcopy(message) for message in self._stored.messages] - return _reconstruct_messages_from_thread_snapshot( - stored_messages=self._stored.messages, - incoming_messages=incoming, - stored_interrupt=self._stored.interrupt, - ) + return [copy.deepcopy(message) for message in self._stored.messages] + incoming async def save( self, diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py index 7580ca77852..23b9b2341e4 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_workflow.py @@ -446,9 +446,16 @@ async def run(self, input_data: dict[str, Any]) -> AsyncGenerator[BaseEvent]: run_checkpoint_storage = _OwnedWorkflowCheckpointStorage(checkpoint_storage, request_owner) builder_seed_messages = raw_messages if resume_payload is not None or (checkpoint_id is not None and not raw_messages): - # Resume / checkpoint-only requests need stored history. Session-owned - # resume_seeded_messages covers empty vs client-replayed transcripts (#8140). - builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) + # Resume / checkpoint-only requests need stored history. Empty input prepends; + # non-empty/replayed transcripts overlap-merge (#8140). + if builder_seed_messages: + builder_seed_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages if stored_snapshot is not None else [], + incoming_messages=builder_seed_messages, + stored_interrupt=stored_snapshot.interrupt if stored_snapshot is not None else None, + ) + else: + builder_seed_messages = snapshot_session.resume_seeded_messages(builder_seed_messages) snapshot_builder = _WorkflowSnapshotBuilder(builder_seed_messages) if snapshot_session.enabled else None if snapshot_builder is not None and effective_state: # Seed builder state so a run that emits no StateSnapshotEvent still diff --git a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py b/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py deleted file mode 100644 index 5ac120cf2ba..00000000000 --- a/python/packages/ag-ui/tests/ag_ui/test_resume_transcript_dedupe.py +++ /dev/null @@ -1,92 +0,0 @@ -# Copyright (c) Microsoft. All rights reserved. - -"""Regression for AG-UI resume with client-replayed transcript (#8140).""" - -from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot -from agent_framework_ag_ui._snapshot_session import ThreadSnapshotSession -from agent_framework_ag_ui._snapshots import AGUIThreadSnapshot - - -def _stored_history() -> list[dict]: - return [ - {"id": "u1", "role": "user", "content": "please run the tool"}, - { - "id": "a1", - "role": "assistant", - "content": "", - "toolCalls": [{"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}}], - }, - ] - - -def test_resume_with_replayed_transcript_does_not_duplicate_history(): - """Clients that re-send the full transcript on resume must not double-write.""" - stored = _stored_history() - incoming = [ - *stored, - {"id": "u2", "role": "user", "content": "approved"}, - ] - - reconstructed = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored, - incoming_messages=incoming, - stored_interrupt=[{"interruptId": "int-1"}], - ) - - assert [m.get("id") for m in reconstructed] == ["u1", "a1", "u2"] - assert len(reconstructed) == 3 - - -def test_resume_with_partial_new_turn_still_seeds_history(): - """Non-empty resume suffix is overlapped onto stored history without duplication.""" - stored = _stored_history() - incoming = [{"id": "u2", "role": "user", "content": "approved"}] - - reconstructed = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored, - incoming_messages=incoming, - stored_interrupt=[{"interruptId": "int-1"}], - ) - - assert [m.get("id") for m in reconstructed] == ["u1", "a1", "u2"] - - -def test_reconstruct_with_empty_incoming_returns_empty(): - """Empty incoming is unchanged at the helper; ThreadSnapshotSession.resume_seeded_messages seeds instead.""" - stored = _stored_history() - - assert ( - _reconstruct_messages_from_thread_snapshot( - stored_messages=stored, - incoming_messages=[], - stored_interrupt=[{"interruptId": "int-1"}], - ) - == [] - ) - - -def test_resume_seeded_messages_empty_vs_replayed(): - """Session-owned resume seeding covers interrupt-only and replayed transcript shapes.""" - stored = _stored_history() - session = ThreadSnapshotSession( - store=None, - scope=None, - thread_id="t1", - stored=AGUIThreadSnapshot( - messages=stored, - state=None, - interrupt=[{"interruptId": "int-1"}], - session_state=None, - ), - ) - - empty = session.resume_seeded_messages([]) - assert [m.get("id") for m in empty] == ["u1", "a1"] - - replayed = session.resume_seeded_messages( - [ - *stored, - {"id": "u2", "role": "user", "content": "approved"}, - ] - ) - assert [m.get("id") for m in replayed] == ["u1", "a1", "u2"] diff --git a/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py b/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py index d0f7522e1a0..f4ad410ba30 100644 --- a/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py +++ b/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py @@ -209,6 +209,35 @@ async def test_seeded_copies_do_not_alias_stored_snapshot(self) -> None: assert session.stored is not None assert session.stored.messages[0]["content"] == "hi" + async def test_replayed_transcript_uses_reconstructor_not_blind_prepend(self) -> None: + """Client-replayed history must not be naively prepended again (#8140).""" + from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot + + stored = [ + {"id": "u1", "role": "user", "content": "please run the tool"}, + { + "id": "a1", + "role": "assistant", + "content": "", + "toolCalls": [ + {"id": "c1", "type": "function", "function": {"name": "needs_approval", "arguments": "{}"}} + ], + }, + ] + incoming = [*stored, {"id": "u2", "role": "user", "content": "approved"}] + reconstructed = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored, + incoming_messages=incoming, + stored_interrupt=[{"interruptId": "int-1"}], + ) + assert [message["id"] for message in reconstructed] == ["u1", "a1", "u2"] + # resume_seeded_messages remains blind prepend for save-time / empty seeds. + snapshot = AGUIThreadSnapshot(messages=stored, interrupt=[{"interruptId": "int-1"}]) + store = await make_store_with("user-1", "t-replay", snapshot) + session = await ThreadSnapshotSession.open(store=store, scope="user-1", thread_id="t-replay") + seeded = session.resume_seeded_messages([{"id": "u2", "role": "user", "content": "approved"}]) + assert [message["id"] for message in seeded] == ["u1", "a1", "u2"] + async def test_without_stored_snapshot_returns_incoming_unchanged(self) -> None: session = await ThreadSnapshotSession.open(store=None, scope=None, thread_id="t1") incoming = [{"id": "m2", "role": "user", "content": "hello"}] From 10827a5dc172f6b33a719877aeecf2a9f3a40068 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Thu, 10 Sep 2026 10:58:52 +0800 Subject: [PATCH 7/8] test(ag-ui): type folded #8140 snapshot replay fixture for mypy Annotate stored/incoming as list[dict[str, Any]] so Test Typing Checks pass. --- .../packages/ag-ui/tests/ag_ui/test_snapshot_session.py | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py b/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py index f4ad410ba30..4a00e653ca4 100644 --- a/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py +++ b/python/packages/ag-ui/tests/ag_ui/test_snapshot_session.py @@ -8,6 +8,8 @@ that interface only; runner integration is covered by the existing suite. """ +from typing import Any + import pytest from ag_ui.core import ( EventType, @@ -213,7 +215,7 @@ async def test_replayed_transcript_uses_reconstructor_not_blind_prepend(self) -> """Client-replayed history must not be naively prepended again (#8140).""" from agent_framework_ag_ui._run_common import _reconstruct_messages_from_thread_snapshot - stored = [ + stored: list[dict[str, Any]] = [ {"id": "u1", "role": "user", "content": "please run the tool"}, { "id": "a1", @@ -224,7 +226,10 @@ async def test_replayed_transcript_uses_reconstructor_not_blind_prepend(self) -> ], }, ] - incoming = [*stored, {"id": "u2", "role": "user", "content": "approved"}] + incoming: list[dict[str, Any]] = [ + *stored, + {"id": "u2", "role": "user", "content": "approved"}, + ] reconstructed = _reconstruct_messages_from_thread_snapshot( stored_messages=stored, incoming_messages=incoming, From aeff82b8239175fd5efbe6e4ab67a5245c3266b3 Mon Sep 17 00:00:00 2001 From: LI <2484593937@qq.com> Date: Thu, 10 Sep 2026 20:14:35 +0800 Subject: [PATCH 8/8] fix(ag-ui): collapse duplicate reconstruct on generic resume Apply eavanvalkenburg suggestion: always overlap-merge via _reconstruct_messages_from_thread_snapshot; only set seeded_resume_from_snapshot when resume carries a non-empty transcript. --- .../ag-ui/agent_framework_ag_ui/_agent_run.py | 16 +++++----------- 1 file changed, 5 insertions(+), 11 deletions(-) diff --git a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py index cc3274ee9ca..10a734a51aa 100644 --- a/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py +++ b/python/packages/ag-ui/agent_framework_ag_ui/_agent_run.py @@ -2508,17 +2508,11 @@ async def run_agent_stream( # must stay empty so synthesized resume tool messages are the only # turn input; history is restored at save when this flag stays false. seeded_resume_from_snapshot = True - raw_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=raw_messages, - stored_interrupt=stored_snapshot.interrupt, - ) - else: - raw_messages = _reconstruct_messages_from_thread_snapshot( - stored_messages=stored_snapshot.messages, - incoming_messages=raw_messages, - stored_interrupt=stored_snapshot.interrupt, - ) + raw_messages = _reconstruct_messages_from_thread_snapshot( + stored_messages=stored_snapshot.messages, + incoming_messages=raw_messages, + stored_interrupt=stored_snapshot.interrupt, + ) else: provider_suffix, snapshot_seed_messages = _split_service_session_input( stored_snapshot_messages=stored_snapshot.messages,