Repository navigation
fix: preserve completed tool outputs when a sibling fails #5332
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
af69e37
39324dd
4e334fb
8493f58
959ebbb
5a2bc69
6c177ae
c727ff5
57a2c71
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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, | ||
| ) | ||
|
|
@@ -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, | ||
|
|
@@ -136,6 +139,7 @@ | |
| NextStepInterruption, | ||
| NextStepRunAgain, | ||
| ProcessedResponse, | ||
| SingleStepResult, | ||
| ) | ||
| from .run_internal.session_persistence import ( | ||
| _session_get_items, | ||
|
|
@@ -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: | ||
|
|
@@ -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, | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
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 AGENTS.md reference: AGENTS.md:L116-L116 Useful? React with 👍 / 👎. |
||
| ) | ||
| ) | ||
| ) | ||
|
|
||
|
|
@@ -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 | ||
|
|
@@ -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( | ||
|
jbeckwith-oai marked this conversation as resolved.
|
||
| all_input_guardrails | ||
| ) and not input_guardrails_triggered(_attempt_input_guardrail_results()) | ||
|
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) | ||
|
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( | ||
|
|
||
There was a problem hiding this comment.
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
Fresh evidence beyond the earlier stream-completion thread: when streaming with parallel input guardrails, a model-selected tool that propagates
CancelledErrorsets_tool_error_selected; this condition then consumes a late tripwire without storingInputGuardrailTripwireTriggered. 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, orother. Example:@codex security dismiss duplicate Already flagged by another reviewWhat each reason means
false-positive— Not a vulnerabilityduplicate— Already tracked elsewhereout-of-scope— Outside this review's scopecompensating-control— Mitigated by another controlrisk-accepted— Risk intentionally acceptedother— Another reason; context requiredUseful? React with 👍 / 👎.