From c65aebccf02a04b8151268f32647d3a3148b08de Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 09:54:04 -0700 Subject: [PATCH] review-fix(conversation-loop): finalize_turn binds stream-recovered final_response before fallible persist sub-steps (BASE order) On BASE the #95514 stream-recovery rebound final_response inline, before the tail-closing / persist-override / micro-compaction / _persist_session calls inside the guarded persist try-block, so a raise in any of them left the caller with the streamed text (cleanup_errors reported the failure). HEAD's _close_transcript_tail returned the recovered value only on normal completion; an exception in the same block dropped it and the user-visible answer regressed to the pre-recovery '' (then rewritten by the explainer). Split the helper into _drop_transcript_scaffolding / _recover_final_from_stream / _close_transcript_tail and rebind final_response in the guarded step as soon as it is computed, restoring BASE ordering. Adds a regression test. --- agent/turn_finalizer.py | 50 ++++++++++++------- ...rn_finalizer_final_response_persistence.py | 38 ++++++++++++++ 2 files changed, 70 insertions(+), 18 deletions(-) diff --git a/agent/turn_finalizer.py b/agent/turn_finalizer.py index 7fcea4d05d..f2ebf62d17 100644 --- a/agent/turn_finalizer.py +++ b/agent/turn_finalizer.py @@ -189,25 +189,32 @@ def _rollback_interrupted_preflight_display(agent, interrupted) -> None: _rollback_fn(_preflight_snapshot) -def _close_transcript_tail(agent, messages, final_response, interrupted, failed): - """Shape the transcript tail before the durable snapshot; returns the (possibly - stream-recovered) ``final_response``.""" - # Strip private retry scaffolding first, or a later "continue" replays - # assistant("(empty)") / recovery nudges into the same empty-response loop. Only - # the synthetic verification nudges go; the assistant candidate persists (#65919). +def _drop_transcript_scaffolding(agent, messages) -> None: + """Strip private retry scaffolding first, or a later "continue" replays + assistant("(empty)") / recovery nudges into the same empty-response loop. Only + the synthetic verification nudges go; the assistant candidate persists (#65919).""" agent._drop_trailing_empty_response_scaffolding(messages) _drop_verification_continuation_scaffolding(messages) - # An empty terminal completion is not authoritative when the stream already - # delivered text; recover before persist so a blank tail isn't frozen (#95514). - _recovered_from_stream = False - if not interrupted and not failed: - _streamed = getattr(agent, "_current_streamed_assistant_text", "") or "" - _streamed = _streamed.strip() if isinstance(_streamed, str) else "" - if not (flatten_message_text(final_response).strip() if final_response else "") and _streamed: - final_response = _streamed - _recovered_from_stream = True +def _recover_final_from_stream(agent, final_response, interrupted, failed) -> Tuple[Any, bool]: + """An empty terminal completion is not authoritative when the stream already + delivered text; recover before persist so a blank tail isn't frozen (#95514). + Returns ``(final_response, recovered_from_stream)``. Called by the finalizer BEFORE + the fallible tail-shaping/persist steps so the recovered text is already bound when + one of them raises — a persist failure must not lose text the user already saw.""" + if interrupted or failed: + return final_response, False + _streamed = getattr(agent, "_current_streamed_assistant_text", "") or "" + _streamed = _streamed.strip() if isinstance(_streamed, str) else "" + if not (flatten_message_text(final_response).strip() if final_response else "") and _streamed: + return _streamed, True + return final_response, False + + +def _close_transcript_tail(agent, messages, final_response, interrupted, _recovered_from_stream) -> None: + """Shape the transcript tail before the durable snapshot (scaffolding already dropped + and ``final_response`` already stream-recovered by the caller).""" # An interrupt can leave a tool result as the tail; close the sequence so strict # providers don't see ``tool → user`` (placeholder: final_response is usually empty). if interrupted: @@ -250,7 +257,6 @@ def _close_transcript_tail(agent, messages, final_response, interrupted, failed) _apply_override = getattr(agent, "_apply_persist_user_message_override", None) if callable(_apply_override): _apply_override(messages) - return final_response def _micro_compact_after_turn(agent, messages, final_response, logger) -> None: @@ -459,10 +465,18 @@ def finalize_turn( "cleanup_task_resources", lambda: agent._cleanup_task_resources(effective_task_id), _cleanup_errors, logger, ) - # Persist only after the transcript tail is shaped and scaffolding removed. + # Persist only after the transcript tail is shaped and scaffolding removed. Each + # sub-step runs in the same order as the original inline block, and the + # stream-recovered ``final_response`` is rebound the moment it is computed — BEFORE + # the fallible tail-shaping / override / micro-compaction / persist calls — so a + # raise in any of them can't drop text the user already saw (#95514, #8049). def _persist_step(): nonlocal final_response - final_response = _close_transcript_tail(agent, messages, final_response, interrupted, failed) + _drop_transcript_scaffolding(agent, messages) + final_response, _recovered_from_stream = _recover_final_from_stream( + agent, final_response, interrupted, failed + ) + _close_transcript_tail(agent, messages, final_response, interrupted, _recovered_from_stream) if not interrupted and not failed: _micro_compact_after_turn(agent, messages, final_response, logger) agent._persist_session(messages, conversation_history) diff --git a/tests/agent/test_turn_finalizer_final_response_persistence.py b/tests/agent/test_turn_finalizer_final_response_persistence.py index e674c85e3a..31fe400ea5 100644 --- a/tests/agent/test_turn_finalizer_final_response_persistence.py +++ b/tests/agent/test_turn_finalizer_final_response_persistence.py @@ -366,6 +366,44 @@ def test_failed_turn_does_not_recover_stream_buffer_as_final_response(monkeypatc assert result["failed"] is True +def test_stream_recovered_final_response_survives_persist_step_failure(monkeypatch): + """#95514 + #8049 ordering: the stream-recovered ``final_response`` is bound as soon + as it is computed, BEFORE the fallible tail-shaping / override / persist calls in the + guarded persist step. A raise later in that step (here the persist-override) must not + drop text the user already saw — the caller still gets the streamed answer and the + failure is reported via ``cleanup_errors``.""" + monkeypatch.setattr("hermes_cli.plugins.invoke_hook", lambda *_a, **_kw: []) + agent = FakeAgent() + agent._current_streamed_assistant_text = " streamed answer " + + def exploding_override(_messages): + raise RuntimeError("override exploded") + + agent._apply_persist_user_message_override = exploding_override + messages = [{"role": "user", "content": "q"}, {"role": "assistant", "content": ""}] + + result = finalize_turn( + agent, + final_response="", + api_call_count=2, + interrupted=False, + failed=False, + messages=messages, + conversation_history=[], + effective_task_id="task", + turn_id="turn", + user_message="q", + original_user_message="q", + _should_review_memory=False, + _turn_exit_reason="text_response(final)", + ) + + assert result["final_response"] == "streamed answer" + assert result["cleanup_errors"] == ["persist_session: override exploded"] + # The blank tail was filled before the raise (same order as the inline BASE block). + assert messages[-1]["content"] == "streamed answer" + + def test_delivery_only_reasoning_excerpt_does_not_fill_blank_assistant(monkeypatch): """Labeled empty-terminal excerpt is delivery-only, not durable content.