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
This commit is contained in:
@@ -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:
|
||||
|
||||
90
tests/agent/test_codex_stream_supersession.py
Normal file
90
tests/agent/test_codex_stream_supersession.py
Normal file
@@ -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 == []
|
||||
Reference in New Issue
Block a user