From 09f307d3901d9956396bd10cbe00d11cefb0a7b5 Mon Sep 17 00:00:00 2001 From: lbo728 Date: Sat, 19 Sep 2026 00:57:47 -0700 Subject: [PATCH] fix(codex): superseded streams keep consuming so the final text is not truncated When a newer attempt claimed the stream-writer slot, run_codex_stream's accept_chunk hook returned False and Relay stopped consuming the Codex Responses stream. The assembler then returned a "completed" response holding only the text seen so far, and the gateway delivered a reply cut mid-sentence (#69486) even with display streaming disabled. Supersession now fences only the live callbacks (text and reasoning deltas, commentary, first-delta) through the attempt's writer token and consumption continues to the terminal frame, so the assembled final response is complete. The single-writer invariant still holds: a superseded attempt never writes into the live display. Request retirement (watchdog kill) is unchanged and still stops the worker. Ported from PR #69502 onto the Relay-backed stream path. Fixes #69486 --- agent/codex_runtime.py | 42 +++++---- tests/agent/test_codex_stream_supersession.py | 90 +++++++++++++++++++ 2 files changed, 116 insertions(+), 16 deletions(-) create mode 100644 tests/agent/test_codex_stream_supersession.py diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 7f9f4dba2f..ee611eebfc 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -862,19 +862,37 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta if watchdog_state is not None else getattr(agent, "_active_codex_stream_request_token", None) ) - # Delta-sink claim for the CURRENT physical attempt (None until the stream opens). - writer_token = {"value": None} + # Delta-sink claim for the CURRENT physical attempt (None until the stream opens). A newer attempt that + # claims the sink supersedes this token; that only silences OUR live callbacks — consumption continues, + # because stopping here handed the gateway a "completed" response missing its tail (#69486). + writer_token = {"value": None, "superseded_logged": False} def _request_is_current() -> bool: return request_token is None or getattr(agent, "_active_codex_stream_request_token", None) is request_token + def _writer_is_current() -> bool: + token = writer_token["value"] + if token is None or stream_writer_is_current(agent, token): + return True + if not writer_token["superseded_logged"]: + writer_token["superseded_logged"] = True + logger.warning("Codex streaming attempt superseded by a newer stream; suppressing its live deltas while " + "consuming to completion so the final response is not truncated (model=%s).", + api_kwargs.get("model", "unknown")) + return False + def _fenced(fn: Callable[[Any], None]) -> Callable[[Any], None]: """Wrap a callback so a retired request's late frames never reach the agent.""" return lambda value: fn(value) if _request_is_current() else None + def _live(fn: Callable[..., None]) -> Callable[..., None]: + """Wrap a live-display callback so a superseded writer's frames never reach the sink (retired ones neither).""" + return lambda *args: fn(*args) if _request_is_current() and _writer_is_current() else None + def _on_text_delta(text: str) -> None: agent._codex_streamed_text_parts.append(text) - agent._fire_stream_delta(text) + if _writer_is_current(): + agent._fire_stream_delta(text) def _on_event(event: Any) -> None: # TTFB/activity touch — once per SSE event. now = time.time() @@ -920,14 +938,6 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta # Claim the delta sink for THIS attempt; a newer attempt supersedes this token. writer_token["value"] = claim_stream_writer(agent) - def _accept_codex_chunk(_chunk: Any) -> bool: - token = writer_token["value"] - if token is None or stream_writer_is_current(agent, token): - return True - logger.warning("Codex streaming attempt superseded by a newer stream; stopping consumption to preserve " - "the single-writer invariant (model=%s).", api_kwargs.get("model", "unknown")) - return False - def _drain_for_finalizer(event_stream: Any) -> None: # ``final`` is already assembled; draining only lets Relay run its finalizer. A transport error # here must NOT discard the completed, already-billed response. @@ -954,7 +964,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta agent._abort_request_openai_client(active_client, reason="codex_stream_close_failed") show_commentary = getattr(agent, "show_commentary", True) wants_commentary = getattr(agent, "interim_assistant_callback", None) is not None and show_commentary - on_commentary_message = _fenced(lambda text: agent._fire_streamed_codex_commentary(text)) if wants_commentary else None + on_commentary_message = _live(agent._fire_streamed_codex_commentary) if wants_commentary else None call_role = ("delegated" if getattr(agent, "is_subagent", False) else "fallback" if int(getattr(agent, "_fallback_index", 0) or 0) > 0 else "primary") for attempt in range(max_stream_retries + 1): @@ -968,7 +978,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta with watchdog_state.lock: watchdog_state.retry_started_ts = time.time() intercepted_events: list = [] - writer_token["value"] = event_stream = None + writer_token["value"], writer_token["superseded_logged"], event_stream = None, False, None try: try: event_stream = relay_llm.stream( @@ -977,7 +987,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta name=str(getattr(agent, "provider", "") or "codex"), model_name=str(model or ""), finalizer=lambda: _consume_codex_event_stream(list(intercepted_events), model=model), on_stream_created=_codex_stream_created, on_chunk=intercepted_events.append, - chunk_adapter=lambda chunk: chunk, accept_chunk=_accept_codex_chunk, + chunk_adapter=lambda chunk: chunk, completed_response_predicate=lambda r: bool(hasattr(r, "output") and not hasattr(r, "__iter__")), metadata={"api_mode": "codex_responses", "call_role": call_role, "retry_count": attempt, "api_request_id": getattr(agent, "_current_api_request_id", None)}, @@ -985,8 +995,8 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta ) final = _consume_codex_event_stream( event_stream, model=model, on_text_delta=_fenced(_on_text_delta), - on_reasoning_delta=_fenced(lambda text: agent._fire_reasoning_delta(text)), - on_commentary_message=on_commentary_message, on_first_delta=on_first_delta, + on_reasoning_delta=_live(agent._fire_reasoning_delta), on_commentary_message=on_commentary_message, + on_first_delta=_live(on_first_delta) if on_first_delta is not None else None, on_event=_fenced(_on_event), interrupt_check=_interrupt_or_superseded, ) except transport_errors as exc: diff --git a/tests/agent/test_codex_stream_supersession.py b/tests/agent/test_codex_stream_supersession.py new file mode 100644 index 0000000000..af17cc7ab1 --- /dev/null +++ b/tests/agent/test_codex_stream_supersession.py @@ -0,0 +1,90 @@ +"""A Codex Responses stream that loses the delta sink to a newer attempt keeps consuming (#69486). + +Stopping consumption on supersession returned a ``completed`` response missing its tail, so the gateway +delivered a truncated reply. Supersession must fence only the live callbacks (text/reasoning deltas, +commentary, first-delta) and still assemble the complete final response. + +Grafted from PR #69502 (@byungsker) onto the Relay-backed ``run_codex_stream``. +""" + +from types import SimpleNamespace + +from agent.codex_runtime import run_codex_stream + + +class _FakeCodexClient: + def __init__(self, events): + self.responses = SimpleNamespace(create=lambda **kwargs: iter(events)) + + +def _completed(): + return SimpleNamespace(type="response.completed", + response=SimpleNamespace(id="resp_1", status="completed", usage=None, output=[], + incomplete_details=None, error=None)) + + +def _message_added(phase=None, item_id="m1"): + return SimpleNamespace(type="response.output_item.added", + item=SimpleNamespace(type="message", role="assistant", phase=phase, id=item_id)) + + +def _agent(supersede_after_checks: int): + """Duck-typed agent whose writer fence reports supersession after ``supersede_after_checks`` checks.""" + live = {"deltas": [], "reasoning": [], "commentary": [], "first_delta": 0, "checks": 0} + + def is_current(_token): + live["checks"] += 1 + return live["checks"] <= supersede_after_checks + + agent = SimpleNamespace( + _interrupt_requested=False, show_commentary=True, + _claim_stream_writer=lambda: 1, _stream_writer_is_current=is_current, + _fire_stream_delta=live["deltas"].append, _fire_reasoning_delta=live["reasoning"].append, + _fire_streamed_codex_commentary=live["commentary"].append, + interim_assistant_callback=lambda *a, **k: None, + _touch_activity=lambda _message: None, _client_log_context=lambda: "test-context", + ) + return agent, live + + +def test_superseded_stream_assembles_complete_final_and_fences_live_deltas(): + events = [ + _message_added(), + SimpleNamespace(type="response.output_text.delta", delta="I've added the live"), + SimpleNamespace(type="response.output_text.delta", delta=" tail."), + SimpleNamespace(type="response.reasoning_text.delta", delta="late reasoning"), + _completed(), + ] + agent, live = _agent(supersede_after_checks=1) + + final = run_codex_stream(agent, {"model": "gpt-5.6-terra"}, client=_FakeCodexClient(events)) + + assert final.status == "completed" + assert final.output_text == "I've added the live tail." + assert agent._codex_streamed_text_parts == ["I've added the live", " tail."] + assert live["deltas"] == ["I've added the live"] + assert live["reasoning"] == [] + + +def test_superseded_stream_fences_commentary_and_first_delta_but_keeps_final(): + events = [ + _message_added(phase="commentary", item_id="c1"), + SimpleNamespace(type="response.output_text.delta", delta="Let me check.", item_id="c1"), + SimpleNamespace(type="response.output_item.done", + item=SimpleNamespace(type="message", role="assistant", phase="commentary", id="c1", + content=[SimpleNamespace(type="output_text", text="Let me check.")])), + _message_added(item_id="m1"), + SimpleNamespace(type="response.output_text.delta", delta="Done.", item_id="m1"), + _completed(), + ] + agent, live = _agent(supersede_after_checks=0) + first_delta = [] + + final = run_codex_stream(agent, {"model": "gpt-5.6-terra"}, client=_FakeCodexClient(events), + on_first_delta=lambda: first_delta.append(True)) + + assert final.status == "completed" + assert final.output_text == "Done." + assert live["commentary"] == [] + assert live["deltas"] == [] + assert first_delta == []