Skip to content
Open
11 changes: 11 additions & 0 deletions src/agents/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@

from .util._pretty_print import pretty_print_run_error_details

_TOOL_LOCAL_CANCELLATION_ATTR = "_agents_tool_local_cancellation"
_DRAIN_STREAM_EVENTS_ATTR = "_agents_drain_queued_stream_events"
_DATA_REDACTED_ATTR = "_agents_data_redacted"
_DATA_REDACTED_ERROR_MESSAGE = "Error details are redacted."
Expand All @@ -41,6 +42,16 @@ class _RedactedExceptionCancellationError(asyncio.CancelledError, Exception):
"""Payload-free cancellation that remains catchable as an Exception."""


def _mark_tool_local_cancellation(error: asyncio.CancelledError) -> None:
setattr(error, _TOOL_LOCAL_CANCELLATION_ATTR, True)


def _is_tool_local_cancellation(error: BaseException) -> bool:
return isinstance(error, asyncio.CancelledError) and bool(
getattr(error, _TOOL_LOCAL_CANCELLATION_ATTR, False)
)


def _mark_error_to_drain_stream_events(error: BaseException) -> None:
setattr(error, _DRAIN_STREAM_EVENTS_ATTR, True)

Expand Down
2 changes: 1 addition & 1 deletion src/agents/items.py
Original file line number Diff line number Diff line change
Expand Up @@ -448,7 +448,7 @@ class ToolCallOutputItem(RunItemBase[Any]):
"""SDK-only custom data attached to this tool output.

This data is not part of ``raw_item`` and is not sent back to the model when the output item is
replayed as input.
replayed as input. On a failed run, unfinished custom-data extraction may leave this unset.
"""

@property
Expand Down
10 changes: 8 additions & 2 deletions src/agents/result.py
Original file line number Diff line number Diff line change
Expand Up @@ -703,6 +703,8 @@ class RunResultStreaming(RunResultBase):
_triggered_input_guardrail_result: InputGuardrailResult | None = field(default=None, repr=False)
_output_guardrails_task: asyncio.Task[Any] | None = field(default=None, repr=False)
_stored_exception: BaseException | None = field(default=None, repr=False)
_tool_error_selected: bool = field(default=False, init=False, repr=False)
"""A selected tool error owns failure reporting while its input verdict settles."""
_cancel_mode: Literal["none", "immediate", "after_turn"] = field(default="none", repr=False)
_last_processed_response: ProcessedResponse | None = field(default=None, repr=False)
"""The last processed model response. This is needed for resuming from interruptions."""
Expand Down Expand Up @@ -1174,7 +1176,7 @@ def _check_errors(self):
# Fetch all the completed guardrail results from the queue and raise if needed
while not self._input_guardrail_queue.empty():
guardrail_result = self._input_guardrail_queue.get_nowait()
if guardrail_result.output.tripwire_triggered:
if guardrail_result.output.tripwire_triggered and not self._tool_error_selected:

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🛡️ Codex Security Review · Automatically triggered

P2 Badge Security: Raise late tripwires after tool-local cancellation

Fresh evidence beyond the earlier stream-completion thread: when streaming with parallel input guardrails, a model-selected tool that propagates CancelledError sets _tool_error_selected; this condition then consumes a late tripwire without storing InputGuardrailTripwireTriggered. Because _check_errors() also ignores a cancelled run-loop task, stream_events() ends normally. A host relying on its documented exception to reject untrusted input can treat the request as a clean completion. Preserve the tripwire whenever the selected tool failure has no reportable stream exception.

SECURITY.md reference: SECURITY.md:L38-L38


Dismiss this finding: Reply with @codex security dismiss <reason> [context]. Codex will resolve this conversation automatically; GitHub may require a page refresh to show the result.

Valid reasons: false-positive, duplicate, out-of-scope, compensating-control, risk-accepted, or other. Example: @codex security dismiss duplicate Already flagged by another review

What each reason means
  • false-positive — Not a vulnerability
  • duplicate — Already tracked elsewhere
  • out-of-scope — Outside this review's scope
  • compensating-control — Mitigated by another control
  • risk-accepted — Risk intentionally accepted
  • other — Another reason; context required

Useful? React with 👍 / 👎.

tripwire_exc = InputGuardrailTripwireTriggered(guardrail_result)
tripwire_exc.run_data = self._create_error_details()
self._stored_exception = tripwire_exc
Expand All @@ -1192,7 +1194,11 @@ def _check_errors(self):
run_impl_exc.run_data = self._create_error_details()
self._stored_exception = run_impl_exc

if self._input_guardrails_task and self._input_guardrails_task.done():
if (
not self._tool_error_selected
and self._input_guardrails_task
and self._input_guardrails_task.done()
):
if not self._input_guardrails_task.cancelled():
in_guard_exc = self._input_guardrails_task.exception()
if isinstance(in_guard_exc, Exception):
Expand Down
138 changes: 103 additions & 35 deletions src/agents/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
_clear_data_redacted_error_traceback,
_detach_data_redacted_error_traceback,
_is_error_data_redacted,
_is_tool_local_cancellation,
_prepare_data_redacted_error,
_raise_data_redacted_error,
)
Expand Down Expand Up @@ -120,9 +121,11 @@
from .run_internal.run_grouping import resolve_run_grouping_id
from .run_internal.run_loop import (
_safe_redacted_persistence_error,
_ToolTaskCancellation,
cleanup_models_after_run,
finalize_max_turns_handler_output,
get_output_schema,
preserve_tool_task_cancellation,
resolve_interrupted_turn,
run_input_guardrails,
run_output_guardrails,
Expand All @@ -136,6 +139,7 @@
NextStepInterruption,
NextStepRunAgain,
ProcessedResponse,
SingleStepResult,
)
from .run_internal.session_persistence import (
_session_get_items,
Expand Down Expand Up @@ -1759,6 +1763,7 @@ async def _save_max_turns_handler_output(
)
if current_turn_span is not None:
current_turn_span.start(mark_as_current=True)
partial_tool_results: list[SingleStepResult] = []
try:
if current_turn <= 1:
try:
Expand All @@ -1785,29 +1790,32 @@ async def _save_max_turns_handler_output(
raise

model_task = asyncio.create_task(
run_single_turn(
bindings=current_bindings,
original_input=original_input,
generated_items=items_for_model,
hooks=hooks,
context_wrapper=context_wrapper,
run_config=run_config,
should_run_agent_start_hooks=should_run_agent_start_hooks,
tool_use_tracker=tool_use_tracker,
server_conversation_tracker=server_conversation_tracker,
session=session,
session_items_to_rewind=(
last_saved_input_snapshot_for_rewind
if not is_resumed_state and session_persistence_enabled
else None
),
reasoning_item_id_policy=resolved_reasoning_item_id_policy,
prompt_cache_key_resolver=prompt_cache_key_resolver,
error_handlers=error_handlers,
agent_span=current_span,
on_response_accepted=_commit_pending_server_response,
on_response_hooks_started=_mark_response_hooks_started,
run_state=run_state,
preserve_tool_task_cancellation(
run_single_turn(
bindings=current_bindings,
original_input=original_input,
generated_items=items_for_model,
hooks=hooks,
context_wrapper=context_wrapper,
run_config=run_config,
should_run_agent_start_hooks=should_run_agent_start_hooks,
tool_use_tracker=tool_use_tracker,
server_conversation_tracker=server_conversation_tracker,
session=session,
session_items_to_rewind=(
last_saved_input_snapshot_for_rewind
if not is_resumed_state and session_persistence_enabled
else None
),
reasoning_item_id_policy=resolved_reasoning_item_id_policy,
prompt_cache_key_resolver=prompt_cache_key_resolver,
error_handlers=error_handlers,
agent_span=current_span,
on_response_accepted=_commit_pending_server_response,
on_response_hooks_started=_mark_response_hooks_started,
run_state=run_state,
on_tool_execution_error=partial_tool_results.append,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve selected tool errors in non-streaming guardrail races

When a parallel input guardrail trips after a tool failure has been selected but while that tool batch is still draining post-invocation work, the non-streaming invocation wires only the partial-output callback and never receives the on_tool_error_selected signal used by the streamed path. The guardrail task can therefore finish first, causing the dedicated tripwire branch to cancel model_task and raise InputGuardrailTripwireTriggered, whereas the same ordering in run_streamed() preserves the selected tool error; this hides the actual tool failure and makes error handling depend on runner mode. Wire the selection signal through this path and apply the same precedence before handling a late tripwire.

AGENTS.md reference: AGENTS.md:L116-L116

Useful? React with 👍 / 👎.

)
)
)

Expand Down Expand Up @@ -1843,18 +1851,23 @@ async def _save_max_turns_handler_output(
)
)
raise
except BaseException:
# A non-tripwire failure (the model turn raising, or a
# guardrail raising a non-tripwire error) propagates from
# gather without cancelling the sibling task. Cancel and drain
# whichever side is still pending so it is not left running
# after the run has failed and its exception is not swallowed.
for pending_task in (guardrail_task, model_task):
if not pending_task.done():
pending_task.cancel()
await asyncio.gather(
guardrail_task, model_task, return_exceptions=True
)
except BaseException as error:
try:
if partial_tool_results and (
not isinstance(error, asyncio.CancelledError)
or _is_tool_local_cancellation(error)
):
# Settle admission without replacing the selected tool
# error. Only successful verdicts admit partial history.
await asyncio.wait((guardrail_task,))
finally:
# Parent cancellation still cancels and drains both tasks.
for pending_task in (guardrail_task, model_task):
if not pending_task.done():
pending_task.cancel()
await asyncio.gather(
guardrail_task, model_task, return_exceptions=True
)
raise
else:
turn_result = await model_task
Expand Down Expand Up @@ -1882,7 +1895,62 @@ async def _save_max_turns_handler_output(
on_response_accepted=_commit_pending_server_response,
on_response_hooks_started=_mark_response_hooks_started,
run_state=run_state,
on_tool_execution_error=partial_tool_results.append,
)
except (Exception, asyncio.CancelledError) as error:
if isinstance(
error, asyncio.CancelledError
) and not _is_tool_local_cancellation(error):
raise
if not partial_tool_results:
if isinstance(error, _ToolTaskCancellation):
raise error.error from None
raise
input_accepted = len(_attempt_input_guardrail_results()) >= len(
Comment thread
jbeckwith-oai marked this conversation as resolved.
all_input_guardrails
) and not input_guardrails_triggered(_attempt_input_guardrail_results())
Comment thread
jbeckwith-oai marked this conversation as resolved.
if partial_tool_results and input_accepted:
partial_result = partial_tool_results[0]
generated_items.extend(partial_result.new_step_items)
session_items.extend(partial_result.new_step_items)
model_responses.append(partial_result.model_response)
Comment thread
jbeckwith-oai marked this conversation as resolved.
tool_input_guardrail_results.extend(
partial_result.tool_input_guardrail_results
)
tool_output_guardrail_results.extend(
partial_result.tool_output_guardrail_results
)
if run_state is not None:
_synchronize_accepted_run_state(
run_state,
generated_items=generated_items,
session_items=session_items,
model_responses=model_responses,
tool_input_guardrail_results=tool_input_guardrail_results,
tool_output_guardrail_results=tool_output_guardrail_results,
current_turn=current_turn,
)
run_state._current_step = NextStepRunAgain()
run_state.set_tool_use_tracker_snapshot(
_tool_use_tracker_snapshot()
)
try:
if session_persistence_enabled:
await save_result_to_session(
session,
[],
partial_result.new_step_items,
run_state,
response_id=partial_result.model_response.response_id,
store=store_setting,
wrapper=context_wrapper,
resumed_write_state=run_state,
)
except Exception:
logger.warning("Failed to save completed tools after a tool error")
if isinstance(error, _ToolTaskCancellation):
raise error.error from None
raise
finally:
if current_turn_span is not None:
attach_usage_to_span(
Expand Down
1 change: 1 addition & 0 deletions src/agents/run_internal/guardrails.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,7 @@ async def run_input_guardrails_with_queue(
isinstance(error, Exception)
and asyncio.current_task() is streamed_result._input_guardrails_task
and not streamed_result.is_complete
and not streamed_result._tool_error_selected
):
if streamed_result.run_loop_task and not streamed_result.run_loop_task.done():
streamed_result.run_loop_task.cancel()
Expand Down
Loading
Loading