Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions run_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -4240,6 +4240,7 @@ def _run_codex_stream(self, api_kwargs: dict, client: Any = None, on_first_delta
collected_output_items: list = []
try:
with active_client.responses.stream(**api_kwargs) as stream:
self._touch_activity("receiving Codex stream response")
for event in stream:
if self._interrupt_requested:
break
Expand All @@ -4249,6 +4250,7 @@ def _run_codex_stream(self, api_kwargs: dict, client: Any = None, on_first_delta
delta_text = getattr(event, "delta", "")
if delta_text:
self._codex_streamed_text_parts.append(delta_text)
self._touch_activity("receiving Codex stream response")
if delta_text and not has_tool_calls:
if not first_delta_fired:
first_delta_fired = True
Expand All @@ -4266,6 +4268,7 @@ def _run_codex_stream(self, api_kwargs: dict, client: Any = None, on_first_delta
reasoning_text = getattr(event, "delta", "")
if reasoning_text:
self._fire_reasoning_delta(reasoning_text)
self._touch_activity("receiving Codex reasoning")
# Collect completed output items — some backends
# (chatgpt.com/backend-api/codex) stream valid items
# via response.output_item.done but the SDK's
Expand Down Expand Up @@ -4368,6 +4371,7 @@ def _run_codex_create_stream_fallback(self, api_kwargs: dict, client: Any = None
event_type = getattr(event, "type", None)
if not event_type and isinstance(event, dict):
event_type = event.get("type")
self._touch_activity("receiving Codex fallback stream response")

# Collect output items and text deltas for backfill
if event_type == "response.output_item.done":
Expand Down
50 changes: 50 additions & 0 deletions tests/run_agent/test_run_agent_codex_responses.py
Original file line number Diff line number Diff line change
Expand Up @@ -1197,3 +1197,53 @@ def test_preflight_codex_input_deduplicates_reasoning_ids(monkeypatch):
assert reasoning_ids.count("rs_xyz") == 1
assert reasoning_ids.count("rs_zzz") == 1
assert len(reasoning_items) == 2


def test_codex_stream_updates_activity_tracker(monkeypatch):
"""Regression for #7794: Codex stream events must update _touch_activity."""
agent = _build_agent(monkeypatch)

text_event = SimpleNamespace(type="response.output_text.delta", delta="hello")
reasoning_event = SimpleNamespace(type="response.reasoning.delta", delta="thinking")

class _StreamWithEvents:
def __enter__(self):
return self
def __exit__(self, *a):
return False
def __iter__(self):
yield text_event
yield reasoning_event
def get_final_response(self):
return _codex_message_response("done")

agent.client = SimpleNamespace(
responses=SimpleNamespace(stream=lambda **kw: _StreamWithEvents())
)

initial_ts = agent._last_activity_ts
agent._run_codex_stream(_codex_request_kwargs())
assert agent._last_activity_ts > initial_ts
assert "Codex" in agent._last_activity_desc


def test_codex_fallback_stream_updates_activity_tracker(monkeypatch):
"""Regression for #7794: Codex fallback stream must also update activity."""
agent = _build_agent(monkeypatch)

events = [
SimpleNamespace(type="response.output_text.delta", delta="hi"),
SimpleNamespace(
type="response.completed",
response=_codex_message_response("result"),
),
]

agent.client = SimpleNamespace(
responses=SimpleNamespace(create=lambda **kw: iter(events))
)

initial_ts = agent._last_activity_ts
agent._run_codex_create_stream_fallback(_codex_request_kwargs())
assert agent._last_activity_ts > initial_ts
assert "Codex" in agent._last_activity_desc