mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-21 16:18:55 +00:00
fix(codex): claim the stream-writer token on the codex_responses path too
Widen the #65991 single-writer fence to run_codex_stream: each codex attempt claims the delta sink before consuming events, and the consume loop's interrupt_check now also stops the instant a newer attempt supersedes this one. Parity with the chat_completions / anthropic / bedrock paths from the salvaged fix. Two regression tests: superseded codex stream is fenced mid-stream; sole-writer codex stream delivers unchanged.
This commit is contained in:
parent
35cbffd5c8
commit
d32a6d4cca
2 changed files with 93 additions and 4 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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"]))
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue