diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 8312cdb53aaa..901638ffb072 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -1143,9 +1143,6 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta agent._codex_stream_last_event_ts = time.time() agent._touch_activity("receiving stream response") - def _interrupt_check() -> bool: - return bool(agent._interrupt_requested) - for attempt in range(max_stream_retries + 1): if agent._interrupt_requested: raise InterruptedError("Agent interrupted before Codex stream retry") @@ -1165,6 +1162,27 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta continue raise + # Claim the delta sink for THIS attempt (#65991) — parity with the + # chat_completions/anthropic/bedrock paths. If a prior attempt's + # stream is somehow still alive, this claim supersedes it so its + # late deltas are fenced out of the turn; conversely, a newer + # attempt supersedes us and the interrupt_check below stops our + # consumption immediately. + _writer_token = agent._claim_stream_writer() + + def _interrupt_or_superseded(_tok=_writer_token) -> bool: + if agent._interrupt_requested: + return True + if not agent._stream_writer_is_current(_tok): + 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 True + return False + try: # Compatibility: some mocks/providers return a concrete response # instead of an iterable. Pass it straight through. @@ -1187,7 +1205,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta ), on_first_delta=on_first_delta, on_event=_on_event, - interrupt_check=_interrupt_check, + interrupt_check=_interrupt_or_superseded, ) except (_httpx.RemoteProtocolError, _httpx.ReadTimeout, _httpx.ConnectError, ConnectionError) as exc: if attempt < max_stream_retries: diff --git a/tests/run_agent/test_stream_single_writer_65991.py b/tests/run_agent/test_stream_single_writer_65991.py index 196fedcfd1d3..b1972f969f7e 100644 --- a/tests/run_agent/test_stream_single_writer_65991.py +++ b/tests/run_agent/test_stream_single_writer_65991.py @@ -138,5 +138,76 @@ class TestSingleWriterLoop: assert "-stale-tail" not in "".join(delivered) +class TestCodexSingleWriter: + """The codex_responses path claims the sink and stops when superseded, + matching the chat_completions/anthropic/bedrock parity added in salvage.""" + + def _codex_event(self, event_type, **fields): + return SimpleNamespace(type=event_type, **fields) + + def test_codex_stream_claims_writer_and_stops_when_superseded(self): + from agent.codex_runtime import run_codex_stream + + agent = _make_agent() + agent.api_mode = "codex_responses" + delivered = [] + agent.stream_delta_callback = lambda t: delivered.append(t) + agent._stream_callback = None + + def event_gen(): + yield self._codex_event( + "response.output_text.delta", delta="first", item_id="i1", + ) + # A concurrent retry supersedes this stream between events. + agent._claim_stream_writer() + yield self._codex_event( + "response.output_text.delta", delta="-stale-tail", item_id="i1", + ) + yield self._codex_event( + "response.completed", + response=SimpleNamespace( + id="r1", status="completed", output=[], usage=None, + ), + ) + + mock_client = MagicMock() + mock_client.responses.create.return_value = event_gen() + + run_codex_stream(agent, {"model": "gpt-5.3-codex"}, client=mock_client) + + assert "".join(delivered) == "first" + assert "-stale-tail" not in "".join(delivered) + + def test_codex_stream_undisturbed_when_sole_writer(self): + from agent.codex_runtime import run_codex_stream + + agent = _make_agent() + agent.api_mode = "codex_responses" + delivered = [] + agent.stream_delta_callback = lambda t: delivered.append(t) + agent._stream_callback = None + + def event_gen(): + yield self._codex_event( + "response.output_text.delta", delta="hello ", item_id="i1", + ) + yield self._codex_event( + "response.output_text.delta", delta="world", item_id="i1", + ) + yield self._codex_event( + "response.completed", + response=SimpleNamespace( + id="r1", status="completed", output=[], usage=None, + ), + ) + + mock_client = MagicMock() + mock_client.responses.create.return_value = event_gen() + + run_codex_stream(agent, {"model": "gpt-5.3-codex"}, client=mock_client) + + assert "".join(delivered) == "hello world" + + if __name__ == "__main__": raise SystemExit(pytest.main([__file__, "-q"]))