from __future__ import annotations import asyncio import contextlib import warnings from typing import TYPE_CHECKING, Any, cast from typing_extensions import Unpack from . import _debug from .agent import Agent from .agent_tool_state import set_agent_tool_state_scope from .exceptions import ( AgentsException, InputGuardrailTripwireTriggered, MaxTurnsExceeded, ModelBehaviorError, OutputGuardrailTripwireTriggered, RunErrorDetails, UserError, _await_data_redacted_error_boundary, _clear_data_redacted_error_traceback, _detach_data_redacted_error_traceback, _is_error_data_redacted, _prepare_data_redacted_error, _raise_data_redacted_error, ) from .guardrail import ( InputGuardrailResult, OutputGuardrailResult, ) from .items import ( InputItem, ItemHelpers, ModelResponse, RunItem, TResponseInputItem, ) from .lifecycle import RunHooks from .logger import log_model_and_tool_action_warning, log_tool_action_warning, logger from .memory import Session from .result import RunResult, RunResultStreaming from .run_config import ( DEFAULT_MAX_TURNS, CallModelData, CallModelInputFilter, ModelInputData, OutputGuardrailBlockedMessageArgs, OutputGuardrailBlockedMessageFormatter, ReasoningItemIdPolicy, RunConfig, RunOptions, ToolErrorFormatter, ToolErrorFormatterArgs, ToolExecutionConfig, ToolNameCollisionPolicy as ToolNameCollisionPolicy, ToolNotFoundBehavior, _coerce_run_config, ) from .run_context import RunContextWrapper, TContext from .run_error_handlers import RunErrorHandlers from .run_internal.agent_bindings import bind_public_agent from .run_internal.agent_runner_helpers import ( append_model_response_if_new, apply_resumed_conversation_settings, attach_usage_to_span, build_interruption_result, build_resumed_stream_debug_extra, ensure_context_wrapper, finalize_conversation_tracking, get_unsent_tool_call_ids_for_interrupted_state, input_guardrails_triggered, resolve_processed_response, resolve_resumed_context, resolve_trace_settings, save_final_turn_items_after_guardrails, save_turn_items_if_needed, should_cancel_parallel_model_task_on_input_guardrail_trip, snapshot_usage, update_run_state_for_interruption, usage_delta, validate_output_guardrails_with_server_managed_conversation, validate_session_conversation_settings, ) from .run_internal.approvals import approvals_from_step from .run_internal.blocked_output import ( OUTPUT_GUARDRAIL_BLOCKED_TOOL_OUTPUT, _blocked_output_failure_items, _BlockedOutputOwnerStarts, _current_response_boundary, _final_turn_items_for_persistence, _has_output_guardrails, _is_terminal_tool_output_response, _resolve_output_guardrail_blocked_message, _retained_items_for_blocked_response, _sanitize_blocked_output_guardrail_results, _should_defer_interrupted_session_items, _synchronize_accepted_run_state, _validate_resumed_session_output_guardrail_safety, ) from .run_internal.error_handlers import ( attach_generic_agent_error, build_run_error_data, resolve_run_error_handler_result, ) from .run_internal.items import ( copy_input_items, normalize_resumed_input, reconcile_nested_history_owned_input_after_rewrite, ) from .run_internal.oai_conversation import OpenAIServerConversationTracker from .run_internal.prompt_cache_key import PromptCacheKeyResolver from .run_internal.run_grouping import resolve_run_grouping_id from .run_internal.run_loop import ( _safe_redacted_persistence_error, cleanup_models_after_run, finalize_max_turns_handler_output, get_all_tools, get_output_schema, initialize_computer_tools, resolve_interrupted_turn, run_input_guardrails, run_output_guardrails, run_single_turn, start_streaming, validate_run_hooks, ) from .run_internal.run_steps import ( NextStepFinalOutput, NextStepHandoff, NextStepInterruption, NextStepRunAgain, ProcessedResponse, ) from .run_internal.session_persistence import ( _session_get_items, admit_pending_input, commit_server_pending_input, persist_session_items_for_guardrail_trip, prepare_input_with_session, reconcile_nested_history_owned_session_item_refs, resume_pending_session_write, resumed_turn_items, save_result_to_session, save_resumed_turn_items, session_items_for_turn, update_run_state_after_resume, ) from .run_internal.tool_use_tracker import ( AgentToolUseTracker, hydrate_tool_use_tracker, serialize_tool_use_tracker, ) from .run_state import RunState from .sandbox.memory.rollouts import terminal_metadata_for_exception from .sandbox.runtime import SandboxRuntime from .tool import dispose_resolved_computers from .tool_guardrails import ToolInputGuardrailResult, ToolOutputGuardrailResult from .tracing import Span, SpanError, agent_span, get_current_trace, task_span, turn_span from .tracing.config import include_task_and_turn_spans from .tracing.context import TraceCtxManager, create_trace_for_run from .tracing.span_data import AgentSpanData, TaskSpanData from .util import _error_tracing DEFAULT_AGENT_RUNNER: AgentRunner = None # type: ignore # the value is set at the end of the module __all__ = [ "AgentRunner", "Runner", "RunConfig", "RunOptions", "RunState", "RunContextWrapper", "ModelInputData", "CallModelData", "CallModelInputFilter", "OutputGuardrailBlockedMessageArgs", "OutputGuardrailBlockedMessageFormatter", "ToolNameCollisionPolicy", "ReasoningItemIdPolicy", "ToolExecutionConfig", "ToolErrorFormatter", "ToolErrorFormatterArgs", "ToolNotFoundBehavior", "DEFAULT_MAX_TURNS", "set_default_agent_runner", "get_default_agent_runner", ] def set_default_agent_runner(runner: AgentRunner | None) -> None: """ WARNING: this class is experimental and not part of the public API It should not be used directly. """ global DEFAULT_AGENT_RUNNER DEFAULT_AGENT_RUNNER = runner if runner is not None else AgentRunner() def _data_redacted_sync_cancellation_source(error: BaseException) -> BaseException | None: """Return the marked task cancellation wrapped by Python 3.10, if any.""" if not issubclass(type(error), asyncio.CancelledError): return None try: context = cast(Any, BaseException.__context__).__get__(error, type(error)) except BaseException: return None if ( context is not None and issubclass(type(context), asyncio.CancelledError) and _is_error_data_redacted(context) ): return cast(BaseException, context) return None def get_default_agent_runner() -> AgentRunner: """ WARNING: this class is experimental and not part of the public API It should not be used directly. """ global DEFAULT_AGENT_RUNNER return DEFAULT_AGENT_RUNNER def _sandbox_memory_rollout_id( *, run_config: RunConfig, conversation_id: str | None, session: Session | None, ) -> str | None: if run_config.sandbox is None: return None return resolve_run_grouping_id( conversation_id=conversation_id, session=session, group_id=run_config.group_id, ) def _sandbox_memory_input( *, memory_input_items_for_persistence: list[TResponseInputItem] | None, original_user_input: str | list[TResponseInputItem] | None, original_input: str | list[TResponseInputItem], ) -> str | list[TResponseInputItem]: if memory_input_items_for_persistence is not None: return list(memory_input_items_for_persistence) if original_user_input is not None: return copy_input_items(original_user_input) return copy_input_items(original_input) class Runner: @classmethod async def run( cls, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], *, context: TContext | None = None, max_turns: int | None = DEFAULT_MAX_TURNS, hooks: RunHooks[TContext] | None = None, run_config: RunConfig | dict[str, Any] | None = None, error_handlers: RunErrorHandlers[TContext] | None = None, previous_response_id: str | None = None, auto_previous_response_id: bool = False, conversation_id: str | None = None, session: Session | None = None, ) -> RunResult: """ Run a workflow starting at the given agent. The agent will run in a loop until a final output is generated. The loop runs like so: 1. The agent is invoked with the given input. 2. If there is a final output (i.e. the agent produces something of type `agent.output_type`), the loop terminates. 3. If there's a handoff, we run the loop again, with the new agent. 4. Else, we run tool calls (if any), and re-run the loop. In two cases, the agent may raise an exception: 1. If the max_turns is exceeded, a MaxTurnsExceeded exception is raised unless handled. 2. If a guardrail tripwire is triggered, the matching tripwire exception is raised, e.g. InputGuardrailTripwireTriggered or OutputGuardrailTripwireTriggered. Note: Only the first agent's input guardrails are run. Args: starting_agent: The starting agent to run. input: The initial input to the agent. You can pass a single string for a user message, or a list of input items. context: The context to run the agent with. max_turns: The maximum number of turns to run the agent for. A turn is defined as one AI invocation (including any tool calls that might occur). Pass ``None`` to disable the turn limit. hooks: An object that receives callbacks on various lifecycle events. run_config: Global settings for the entire agent run. error_handlers: Error handlers keyed by error kind. previous_response_id: The ID of the previous response. If using OpenAI models via the Responses API, this allows you to skip passing in input from the previous turn. auto_previous_response_id: If True, enable Responses API response chaining automatically for the first turn even when no ``previous_response_id`` is supplied yet. conversation_id: The conversation ID (https://platform.openai.com/docs/guides/conversation-state?api-mode=responses). If provided, the conversation will be used to read and write items. Every agent will have access to the conversation history so far, and its output items will be written to the conversation. We recommend only using this if you are exclusively using OpenAI models; other model providers don't write to the Conversation object, so you'll end up having partial conversations stored. session: A session for automatic conversation history management. Returns: A run result containing all the inputs, guardrail results and the output of the last agent. Agents may perform handoffs, so we don't know the specific type of the output. """ runner = DEFAULT_AGENT_RUNNER redacted_error: BaseException | None = None try: return await runner.run( starting_agent, input, context=context, max_turns=max_turns, hooks=hooks, run_config=run_config, error_handlers=error_handlers, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, conversation_id=conversation_id, session=session, ) except BaseException as error: if not _is_error_data_redacted(error): raise _detach_data_redacted_error_traceback(error) redacted_error = error starting_agent = cast(Any, None) input = cast(Any, None) context = cast(Any, None) hooks = cast(Any, None) run_config = cast(Any, None) error_handlers = cast(Any, None) previous_response_id = None auto_previous_response_id = cast(Any, None) conversation_id = None session = cast(Any, None) runner = cast(Any, None) assert redacted_error is not None raise redacted_error from None @classmethod def run_sync( cls, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], *, context: TContext | None = None, max_turns: int | None = DEFAULT_MAX_TURNS, hooks: RunHooks[TContext] | None = None, run_config: RunConfig | dict[str, Any] | None = None, error_handlers: RunErrorHandlers[TContext] | None = None, previous_response_id: str | None = None, auto_previous_response_id: bool = False, conversation_id: str | None = None, session: Session | None = None, ) -> RunResult: """ Run a workflow synchronously, starting at the given agent. Note: This just wraps the `run` method, so it will not work if there's already an event loop (e.g. inside an async function, or in a Jupyter notebook or async context like FastAPI). For those cases, use the `run` method instead. The agent will run in a loop until a final output is generated. The loop runs: 1. The agent is invoked with the given input. 2. If there is a final output (i.e. the agent produces something of type `agent.output_type`), the loop terminates. 3. If there's a handoff, we run the loop again, with the new agent. 4. Else, we run tool calls (if any), and re-run the loop. In two cases, the agent may raise an exception: 1. If the max_turns is exceeded, a MaxTurnsExceeded exception is raised unless handled. 2. If a guardrail tripwire is triggered, the matching tripwire exception is raised, e.g. InputGuardrailTripwireTriggered or OutputGuardrailTripwireTriggered. Note: Only the first agent's input guardrails are run. Args: starting_agent: The starting agent to run. input: The initial input to the agent. You can pass a single string for a user message, or a list of input items. context: The context to run the agent with. max_turns: The maximum number of turns to run the agent for. A turn is defined as one AI invocation (including any tool calls that might occur). Pass ``None`` to disable the turn limit. hooks: An object that receives callbacks on various lifecycle events. run_config: Global settings for the entire agent run. error_handlers: Error handlers keyed by error kind. previous_response_id: The ID of the previous response, if using OpenAI models via the Responses API, this allows you to skip passing in input from the previous turn. auto_previous_response_id: If True, enable Responses API response chaining automatically for the first turn even when no ``previous_response_id`` is supplied yet. conversation_id: The ID of the stored conversation, if any. session: A session for automatic conversation history management. Returns: A run result containing all the inputs, guardrail results and the output of the last agent. Agents may perform handoffs, so we don't know the specific type of the output. """ runner = DEFAULT_AGENT_RUNNER redacted_error: BaseException | None = None try: return runner.run_sync( starting_agent, input, context=context, max_turns=max_turns, hooks=hooks, run_config=run_config, error_handlers=error_handlers, previous_response_id=previous_response_id, conversation_id=conversation_id, session=session, auto_previous_response_id=auto_previous_response_id, ) except BaseException as error: if not _is_error_data_redacted(error): raise _detach_data_redacted_error_traceback(error) redacted_error = error starting_agent = cast(Any, None) input = cast(Any, None) context = cast(Any, None) hooks = cast(Any, None) run_config = cast(Any, None) error_handlers = cast(Any, None) previous_response_id = None auto_previous_response_id = cast(Any, None) conversation_id = None session = cast(Any, None) runner = cast(Any, None) assert redacted_error is not None raise redacted_error from None @classmethod def run_streamed( cls, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], context: TContext | None = None, max_turns: int | None = DEFAULT_MAX_TURNS, hooks: RunHooks[TContext] | None = None, run_config: RunConfig | dict[str, Any] | None = None, previous_response_id: str | None = None, auto_previous_response_id: bool = False, conversation_id: str | None = None, session: Session | None = None, *, error_handlers: RunErrorHandlers[TContext] | None = None, ) -> RunResultStreaming: """ Run a workflow starting at the given agent in streaming mode. The returned result object contains a method you can use to stream semantic events as they are generated. The agent will run in a loop until a final output is generated. The loop runs like so: 1. The agent is invoked with the given input. 2. If there is a final output (i.e. the agent produces something of type `agent.output_type`), the loop terminates. 3. If there's a handoff, we run the loop again, with the new agent. 4. Else, we run tool calls (if any), and re-run the loop. In two cases, the agent may raise an exception: 1. If the max_turns is exceeded, a MaxTurnsExceeded exception is raised unless handled. 2. If a guardrail tripwire is triggered, the matching tripwire exception is raised, e.g. InputGuardrailTripwireTriggered or OutputGuardrailTripwireTriggered. Note: Only the first agent's input guardrails are run. Args: starting_agent: The starting agent to run. input: The initial input to the agent. You can pass a single string for a user message, or a list of input items. context: The context to run the agent with. max_turns: The maximum number of turns to run the agent for. A turn is defined as one AI invocation (including any tool calls that might occur). Pass ``None`` to disable the turn limit. hooks: An object that receives callbacks on various lifecycle events. run_config: Global settings for the entire agent run. error_handlers: Error handlers keyed by error kind. previous_response_id: The ID of the previous response, if using OpenAI models via the Responses API, this allows you to skip passing in input from the previous turn. auto_previous_response_id: If True, enable Responses API response chaining automatically for the first turn even when no ``previous_response_id`` is supplied yet. conversation_id: The ID of the stored conversation, if any. session: A session for automatic conversation history management. Returns: A result object that contains data about the run, as well as a method to stream events. """ runner = DEFAULT_AGENT_RUNNER return runner.run_streamed( starting_agent, input, context=context, max_turns=max_turns, hooks=hooks, run_config=run_config, error_handlers=error_handlers, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, conversation_id=conversation_id, session=session, ) class AgentRunner: """ WARNING: this class is experimental and not part of the public API It should not be used directly or subclassed. """ async def run( self, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], **kwargs: Unpack[RunOptions[TContext]], ) -> RunResult: redacted_error: BaseException | None = None try: return await self._run_impl(starting_agent, input, **kwargs) except BaseException as error: if not _is_error_data_redacted(error): raise _detach_data_redacted_error_traceback(error) redacted_error = error self = cast(Any, None) starting_agent = cast(Any, None) input = cast(Any, None) cast(dict[str, Any], kwargs).clear() assert redacted_error is not None _detach_data_redacted_error_traceback(redacted_error) raise redacted_error from None async def _run_impl( self, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], **kwargs: Unpack[RunOptions[TContext]], ) -> RunResult: context = kwargs.get("context") max_turns = kwargs.get("max_turns", DEFAULT_MAX_TURNS) hooks = cast(RunHooks[TContext], validate_run_hooks(kwargs.get("hooks"))) run_config = kwargs.get("run_config") error_handlers = kwargs.get("error_handlers") previous_response_id = kwargs.get("previous_response_id") auto_previous_response_id = kwargs.get("auto_previous_response_id", False) conversation_id = kwargs.get("conversation_id") session = kwargs.get("session") run_config = RunConfig() if run_config is None else _coerce_run_config(run_config) is_resumed_state = isinstance(input, RunState) run_state: RunState[TContext] | None = ( cast(RunState[TContext], input) if is_resumed_state else None ) resolved_reasoning_item_id_policy: ReasoningItemIdPolicy | None = ( run_config.reasoning_item_id_policy if run_config.reasoning_item_id_policy is not None else (run_state._reasoning_item_id_policy if run_state is not None else None) ) if run_state is not None: run_state._reasoning_item_id_policy = resolved_reasoning_item_id_policy starting_input = input if not is_resumed_state else None original_user_input: str | list[TResponseInputItem] | None = None session_input_items_for_persistence: list[TResponseInputItem] | None = ( [] if (session is not None and is_resumed_state) else None ) # Track the most recent input batch we persisted so conversation-lock retries can rewind # exactly those items (and not the full history). last_saved_input_snapshot_for_rewind: list[TResponseInputItem] | None = None if is_resumed_state and run_state is not None: ( conversation_id, previous_response_id, auto_previous_response_id, ) = apply_resumed_conversation_settings( run_state=run_state, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) validate_session_conversation_settings( session, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) starting_input = run_state._original_input original_user_input = copy_input_items(run_state._original_input) prepared_input = normalize_resumed_input(original_user_input) context_wrapper = resolve_resumed_context( run_state=run_state, context=context, ) context = context_wrapper.context await resume_pending_session_write(run_state, session, wrapper=context_wrapper) max_turns = run_state._max_turns else: raw_input = cast(str | list[TResponseInputItem], input) original_user_input = raw_input validate_session_conversation_settings( session, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) context_wrapper = ensure_context_wrapper(context) context = context_wrapper.context set_agent_tool_state_scope(context_wrapper, None) server_manages_conversation = ( conversation_id is not None or previous_response_id is not None or auto_previous_response_id ) if server_manages_conversation: prepared_input, _ = await prepare_input_with_session( raw_input, session, run_config.session_input_callback, run_config.session_settings, include_history_in_prepared_input=False, preserve_dropped_new_items=True, reasoning_item_id_policy=resolved_reasoning_item_id_policy, wrapper=context_wrapper, ) original_input_for_state = raw_input session_input_items_for_persistence = [] else: ( prepared_input, session_input_items_for_persistence, ) = await prepare_input_with_session( raw_input, session, run_config.session_input_callback, run_config.session_settings, reasoning_item_id_policy=resolved_reasoning_item_id_policy, wrapper=context_wrapper, ) original_input_for_state = prepared_input # Check whether to enable OpenAI server-managed conversation if ( conversation_id is not None or previous_response_id is not None or auto_previous_response_id ): server_conversation_tracker = OpenAIServerConversationTracker( conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, reasoning_item_id_policy=resolved_reasoning_item_id_policy, ) else: server_conversation_tracker = None session_persistence_enabled = session is not None and server_conversation_tracker is None memory_input_items_for_persistence = ( list(session_input_items_for_persistence) if session_persistence_enabled and session_input_items_for_persistence is not None else None ) if server_conversation_tracker is not None and is_resumed_state and run_state is not None: session_input_items: list[TResponseInputItem] | None = None if session is not None: try: session_input_items = await _session_get_items( session, wrapper=context_wrapper, ) except Exception: session_input_items = None server_conversation_tracker.hydrate_from_state( original_input=run_state._original_input, generated_items=run_state._generated_items, model_responses=run_state._model_responses, session_items=session_input_items, unsent_tool_call_ids=get_unsent_tool_call_ids_for_interrupted_state(run_state), ) tool_use_tracker = AgentToolUseTracker() if is_resumed_state and run_state is not None: hydrate_tool_use_tracker(tool_use_tracker, run_state, starting_agent) ( trace_workflow_name, trace_id, trace_group_id, trace_metadata, trace_config, ) = resolve_trace_settings(run_state=run_state, run_config=run_config) with TraceCtxManager( workflow_name=trace_workflow_name, trace_id=trace_id, group_id=trace_group_id, metadata=trace_metadata, tracing=trace_config, disabled=run_config.tracing_disabled, trace_state=run_state._trace_state if run_state is not None else None, reattach_resumed_trace=is_resumed_state, ): if is_resumed_state and run_state is not None: run_state.set_trace(get_current_trace()) current_turn = run_state._current_turn raw_original_input = run_state._original_input original_input = normalize_resumed_input(raw_original_input) ( original_input, run_state._nested_history_owned_session_item_refs, ) = reconcile_nested_history_owned_input_after_rewrite( raw_original_input, original_input, run_state._nested_history_owned_session_item_refs, ) run_state._original_input = copy_input_items(original_input) # Copy every list adopted from the state: the run appends to these, and # the caller still owns the state as a resumable snapshot. generated_items = list(run_state._generated_items) session_items = list(run_state._session_items) model_responses = list(run_state._model_responses) # Cast to the correct type since we know this is TContext context_wrapper = cast(RunContextWrapper[TContext], run_state._context) else: current_turn = 0 original_input = copy_input_items(original_input_for_state) generated_items = [] session_items = [] model_responses = [] run_state = RunState( context=context_wrapper, original_input=original_input, starting_agent=starting_agent, max_turns=max_turns, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) run_state._reasoning_item_id_policy = resolved_reasoning_item_id_policy run_state.set_trace(get_current_trace()) use_task_and_turn_spans = include_task_and_turn_spans(run_config.tracing) current_task_span: Span[TaskSpanData] | None = ( task_span(name=trace_workflow_name) if use_task_and_turn_spans else None ) if current_task_span is not None: current_task_span.start(mark_as_current=True) task_usage_start = snapshot_usage(context_wrapper.usage) try: sandbox_runtime = SandboxRuntime( starting_agent=starting_agent, run_config=run_config, rollout_id=_sandbox_memory_rollout_id( run_config=run_config, conversation_id=conversation_id, session=session, ), run_state=run_state, ) prompt_cache_key_resolver = PromptCacheKeyResolver.from_run_state( run_state=run_state, ) completed_result: RunResult | None = None run_exception: BaseException | None = None def _with_reasoning_item_id_policy(result: RunResult) -> RunResult: result._reasoning_item_id_policy = resolved_reasoning_item_id_policy if run_state is not None: run_state._reasoning_item_id_policy = resolved_reasoning_item_id_policy return result def _tool_use_tracker_snapshot() -> dict[str, list[str]]: identity_root_agent = starting_agent if run_state is not None and run_state._starting_agent is not None: identity_root_agent = run_state._starting_agent return serialize_tool_use_tracker( tool_use_tracker, starting_agent=identity_root_agent, ) def _finalize_result(result: RunResult) -> RunResult: nonlocal completed_result result._starting_agent_for_state = ( run_state._starting_agent if run_state is not None and run_state._starting_agent is not None else starting_agent ) finalized_result = finalize_conversation_tracking( _with_reasoning_item_id_policy(result), server_conversation_tracker=server_conversation_tracker, run_state=run_state, ) sandbox_runtime.apply_result_metadata(finalized_result) if run_state is not None: finalized_result._generated_prompt_cache_key = ( run_state._generated_prompt_cache_key ) finalized_result._pending_input_for_state = run_state.pending_input finalized_result._current_step_for_state = run_state._current_step finalized_result._nested_history_owned_session_item_refs = list( run_state._nested_history_owned_session_item_refs ) completed_result = finalized_result return finalized_result pending_server_items: list[RunItem] | None = None pending_input_admission_items: list[InputItem] = [] input_guardrail_results: list[InputGuardrailResult] = ( list(run_state._input_guardrail_results) if run_state is not None else [] ) input_guardrail_attempt_start = len(input_guardrail_results) def _attempt_input_guardrail_results() -> list[InputGuardrailResult]: return input_guardrail_results[input_guardrail_attempt_start:] def _commit_pending_server_response( model_response: ModelResponse, processed_response: ProcessedResponse | None, ) -> bool: if ( run_state is None or server_conversation_tracker is None or not pending_input_admission_items ): return False return commit_server_pending_input( run_state=run_state, tracker=server_conversation_tracker, admission_items=pending_input_admission_items, generated_items=generated_items, session_items=session_items, model_response=model_response, processed_response=processed_response, current_turn=current_turn, ) def _mark_response_hooks_started() -> None: if run_state is None or not isinstance( run_state._current_step, NextStepInterruption ): return if run_state._current_step.response_accepted: run_state._current_step.llm_end_hooks_started = True # Output guardrails run once, at the end of the run. Accumulate their results # here so the failure handler below can report them on the raised exception. output_guardrail_results: list[OutputGuardrailResult] = ( list(run_state._output_guardrail_results) if run_state is not None else [] ) tool_input_guardrail_results: list[ToolInputGuardrailResult] = ( list(getattr(run_state, "_tool_input_guardrail_results", [])) if run_state is not None else [] ) tool_output_guardrail_results: list[ToolOutputGuardrailResult] = ( list(getattr(run_state, "_tool_output_guardrail_results", [])) if run_state is not None else [] ) current_span: Span[AgentSpanData] | None = None if ( is_resumed_state and run_state is not None and run_state._current_agent is not None ): current_agent = run_state._current_agent else: current_agent = starting_agent _validate_resumed_session_output_guardrail_safety( agent=current_agent, run_config=run_config, session=session, run_state=run_state if is_resumed_state else None, ) sandbox_runtime.assert_agent_supported(current_agent) should_run_agent_start_hooks = True store_setting = current_agent.model_settings.resolve( run_config.model_settings ).store if ( not is_resumed_state and session_persistence_enabled and original_user_input is not None and session_input_items_for_persistence is None ): sandbox_runtime.assert_agent_supported(current_agent) session_input_items_for_persistence = ItemHelpers.input_to_new_input_list( original_user_input ) if ( session_persistence_enabled and session_input_items_for_persistence and not sandbox_runtime.enabled ): # Capture the exact input saved so it can be rewound on conversation # lock retries. last_saved_input_snapshot_for_rewind = list(session_input_items_for_persistence) await save_result_to_session( session, session_input_items_for_persistence, [], run_state, store=store_setting, wrapper=context_wrapper, ) session_input_items_for_persistence = [] except BaseException: if current_task_span is not None: attach_usage_to_span( current_task_span, usage_delta(task_usage_start, context_wrapper.usage), ) current_task_span.finish(reset_current=True) raise try: while True: validate_output_guardrails_with_server_managed_conversation( current_agent, run_config, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) if TYPE_CHECKING: # Keep loop-carried types explicit to bound Pyright's flow analysis. original_input = cast( # type: ignore[redundant-cast] str | list[TResponseInputItem], original_input ) run_state = cast(RunState[TContext] | None, run_state) resuming_turn = is_resumed_state all_input_guardrails = ( starting_agent.input_guardrails + (run_config.input_guardrails or []) if current_turn == 0 and not resuming_turn else [] ) sequential_guardrails = [ g for g in all_input_guardrails if not g.run_in_parallel ] parallel_guardrails = [g for g in all_input_guardrails if g.run_in_parallel] if sandbox_runtime.enabled and sequential_guardrails: # Blocking first-turn guardrails must run before sandbox prep so a tripwire # can prevent session creation, startup, or live-session mutation. try: await run_input_guardrails( starting_agent, sequential_guardrails, copy_input_items(original_input), context_wrapper, input_guardrail_results, ) except InputGuardrailTripwireTriggered: session_input_items_for_persistence = ( await persist_session_items_for_guardrail_trip( session, server_conversation_tracker, session_input_items_for_persistence, original_user_input, run_state, store=store_setting, wrapper=context_wrapper, ) ) raise sequential_guardrails = [] current_bindings = bind_public_agent(current_agent) execution_agent = current_bindings.execution_agent input_before_sandbox = copy_input_items(original_input) prepared_sandbox = await sandbox_runtime.prepare_agent( current_agent=current_agent, current_input=original_input, context_wrapper=context_wrapper, is_resumed_state=resuming_turn, ) current_bindings = prepared_sandbox.bindings execution_agent = current_bindings.execution_agent if run_state is not None: ( original_input, run_state._nested_history_owned_session_item_refs, ) = reconcile_nested_history_owned_input_after_rewrite( input_before_sandbox, prepared_sandbox.input, run_state._nested_history_owned_session_item_refs, ) else: original_input = copy_input_items(prepared_sandbox.input) if starting_input is not None and not isinstance(starting_input, RunState): starting_input = copy_input_items(original_input) if run_state is not None: run_state._original_input = copy_input_items(original_input) normalized_starting_input: str | list[TResponseInputItem] = ( starting_input if starting_input is not None and not isinstance(starting_input, RunState) else "" ) store_setting = current_agent.model_settings.resolve( run_config.model_settings ).store if session_persistence_enabled and session_input_items_for_persistence: last_saved_input_snapshot_for_rewind = list( session_input_items_for_persistence ) await save_result_to_session( session, list(last_saved_input_snapshot_for_rewind), [], run_state, store=store_setting, wrapper=context_wrapper, ) session_input_items_for_persistence = [] if run_state is not None and run_state._current_step is not None: if isinstance(run_state._current_step, NextStepInterruption): logger.debug("Continuing from interruption") if not run_state._model_responses: raise UserError("No model response found in previous state") if run_state._last_processed_response is None: if run_state._current_step.response_accepted: raise UserError( "An accepted model response could not be processed; " "start a new run instead of retrying it" ) raise UserError("No processed response found in previous state") resumed_response_boundary = _current_response_boundary( (), run_state._last_processed_response, run_state, ) blocked_output_owner_starts = _BlockedOutputOwnerStarts( nonstreamed_session_items=(resumed_response_boundary.session_start), run_state_generated_items=( resumed_response_boundary.generated_start ), run_state_session_items=resumed_response_boundary.session_start, run_state_model_responses=len(run_state._model_responses) - 1, run_state_tool_output_guardrail_results=len( run_state._tool_output_guardrail_results ), ) turn_result = await resolve_interrupted_turn( bindings=current_bindings, original_input=original_input, original_pre_step_items=generated_items, new_response=run_state._model_responses[-1], processed_response=run_state._last_processed_response, hooks=hooks, context_wrapper=context_wrapper, run_config=run_config, server_manages_conversation=( server_conversation_tracker is not None ), run_state=run_state, error_handlers=error_handlers, ) if run_state._last_processed_response is not None: tool_use_tracker.record_processed_response( current_agent, run_state._last_processed_response, ) input_before_turn_rewrite = original_input original_input = turn_result.original_input generated_items, turn_session_items = resumed_turn_items(turn_result) session_items.extend(turn_session_items) if run_state is not None: if turn_result.nested_history_owned_items is not None: run_state._nested_history_owned_session_item_refs = ( reconcile_nested_history_owned_session_item_refs( session_items, run_state._nested_history_owned_session_item_refs, input_before_turn_rewrite, turn_result.original_input, turn_result.nested_history_owned_items, ) ) update_run_state_after_resume( run_state, turn_result=turn_result, generated_items=generated_items, session_items=session_items, ) if isinstance( turn_result.next_step, NextStepInterruption | NextStepHandoff, ): # Publish before the fallible append so a retry does not # lose guardrail results for work that already ran. run_state._tool_input_guardrail_results = [ *tool_input_guardrail_results, *turn_result.tool_input_guardrail_results, ] run_state._tool_output_guardrail_results = [ *tool_output_guardrail_results, *turn_result.tool_output_guardrail_results, ] if ( session_persistence_enabled and turn_session_items and run_state is not None and not isinstance(turn_result.next_step, NextStepFinalOutput) and not ( isinstance(turn_result.next_step, NextStepInterruption) and _should_defer_interrupted_session_items( current_agent, run_config, ) ) ): run_state._current_turn_persisted_item_count = ( await save_resumed_turn_items( run_state=run_state, session=session, items=turn_session_items, persisted_count=( run_state._current_turn_persisted_item_count ), response_id=turn_result.model_response.response_id, reasoning_item_id_policy=( run_state._reasoning_item_id_policy ), store=store_setting, wrapper=context_wrapper, ) ) # After the resumed turn, treat subsequent turns as fresh so # counters and input saving behave normally. is_resumed_state = False if isinstance(turn_result.next_step, NextStepInterruption): interruption_result_input: str | list[TResponseInputItem] = ( original_input ) append_model_response_if_new( model_responses, turn_result.model_response ) tool_input_guardrail_results.extend( turn_result.tool_input_guardrail_results ) tool_output_guardrail_results.extend( turn_result.tool_output_guardrail_results ) processed_response_for_state = resolve_processed_response( run_state=run_state, processed_response=turn_result.processed_response, ) if run_state is not None: update_run_state_for_interruption( run_state=run_state, model_responses=model_responses, processed_response=processed_response_for_state, generated_items=generated_items, session_items=session_items, current_turn=current_turn, next_step=turn_result.next_step, ) result = build_interruption_result( result_input=interruption_result_input, session_items=session_items, model_responses=model_responses, current_agent=current_agent, input_guardrail_results=input_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, context_wrapper=context_wrapper, interruptions=approvals_from_step(turn_result.next_step), processed_response=processed_response_for_state, tool_use_tracker=tool_use_tracker, max_turns=max_turns, current_turn=current_turn, generated_items=generated_items, run_state=run_state, original_input=original_input, ) return _finalize_result(result) if isinstance(turn_result.next_step, NextStepRunAgain): continue append_model_response_if_new( model_responses, turn_result.model_response ) tool_input_guardrail_results.extend( turn_result.tool_input_guardrail_results ) tool_output_guardrail_results.extend( turn_result.tool_output_guardrail_results ) if isinstance(turn_result.next_step, NextStepFinalOutput): if run_state is not None and _has_output_guardrails( current_agent, run_config ): run_state._tool_output_guardrail_results = list( tool_output_guardrail_results ) current_processed_response = ( turn_result.processed_response if turn_result.processed_response is not None else run_state._last_processed_response ) output_guardrail_result_start = len(output_guardrail_results) try: await run_output_guardrails( current_agent.output_guardrails + (run_config.output_guardrails or []), current_agent, turn_result.next_step.output, context_wrapper, output_guardrail_results, ) except OutputGuardrailTripwireTriggered as exc: if not _is_terminal_tool_output_response( turn_session_items, current_processed_response, run_state, ): raise sanitized_results = _sanitize_blocked_output_guardrail_results( output_guardrail_results[output_guardrail_result_start:], exc, ) output_guardrail_results[output_guardrail_result_start:] = ( sanitized_results ) session_items = _blocked_output_failure_items( session_items, (), blocked_output_owner_starts, ) blocked_message = _resolve_output_guardrail_blocked_message( exc, agent=current_agent, run_config=run_config, context_wrapper=context_wrapper, ) if blocked_message != OUTPUT_GUARDRAIL_BLOCKED_TOOL_OUTPUT: sanitized_results = ( _sanitize_blocked_output_guardrail_results( sanitized_results, exc, blocked_message, ) ) output_guardrail_results[output_guardrail_result_start:] = ( sanitized_results ) retained_items = _retained_items_for_blocked_response( turn_session_items, turn_result.model_response, run_state, current_processed_response, owner_starts=blocked_output_owner_starts, blocked_message=blocked_message, ) list.extend(session_items, retained_items) try: await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=( session_persistence_enabled ), input_guardrail_results=( _attempt_input_guardrail_results() ), items=retained_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) except BaseException as persistence_error: raise _safe_redacted_persistence_error( persistence_error ) from None raise except (Exception, asyncio.CancelledError) as guardrail_error: if not isinstance( guardrail_error, asyncio.CancelledError ) or not _is_terminal_tool_output_response( turn_session_items, current_processed_response, run_state, ): final_turn_items = _final_turn_items_for_persistence( turn_session_items, current_processed_response, run_state, current_agent, run_config, ) await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=( _attempt_input_guardrail_results() ), items=final_turn_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) raise final_turn_items = _final_turn_items_for_persistence( turn_session_items, current_processed_response, run_state, current_agent, run_config, ) await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=_attempt_input_guardrail_results(), items=final_turn_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) current_step = getattr(run_state, "_current_step", None) approvals_from_state = approvals_from_step(current_step) result = RunResult( input=turn_result.original_input, new_items=session_items, raw_responses=model_responses, final_output=turn_result.next_step.output, _last_agent=current_agent, input_guardrail_results=input_guardrail_results, output_guardrail_results=output_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, context_wrapper=context_wrapper, interruptions=approvals_from_state, _tool_use_tracker_snapshot=_tool_use_tracker_snapshot(), max_turns=max_turns, ) result._current_turn = current_turn result._model_input_items = list(generated_items) # Keep normalized replay aligned with the model-facing # continuation whenever session history preserved extra items. result._replay_from_model_input_items = list( generated_items ) != list(session_items) if run_state is not None: result._trace_state = run_state._trace_state result._original_input = copy_input_items(original_input) run_state._current_step = None return _finalize_result(result) elif isinstance(turn_result.next_step, NextStepHandoff): current_agent = cast( Agent[TContext], turn_result.next_step.new_agent ) if run_state is not None: run_state._current_agent = current_agent starting_input = turn_result.original_input original_input = turn_result.original_input if current_span is not None: current_span.finish(reset_current=True) current_span = None should_run_agent_start_hooks = True continue continue if run_state is not None: if run_state._current_step is None: run_state._current_step = NextStepRunAgain() pending_input = run_state.pending_input if pending_input: pending_guardrails = current_agent.input_guardrails + ( run_config.input_guardrails or [] ) try: await run_input_guardrails( current_agent, pending_guardrails, pending_input, context_wrapper, input_guardrail_results, ) finally: run_state._input_guardrail_results = list(input_guardrail_results) admission_items = await admit_pending_input( run_state=run_state, agent=current_agent, session=session, server_conversation_tracker=server_conversation_tracker, store=store_setting, wrapper=context_wrapper, ) generated_items.extend(admission_items) session_items.extend(admission_items) if pending_server_items is not None: pending_server_items.extend(admission_items) pending_input_admission_items = [ item for item in admission_items if isinstance(item, InputItem) ] if not run_state._pending_input: run_state._generated_items = list(generated_items) run_state._session_items = list(session_items) all_tools = await get_all_tools(execution_agent, context_wrapper) all_tools = await initialize_computer_tools( tools=all_tools, context_wrapper=context_wrapper ) if current_span is None: if (output_schema := get_output_schema(execution_agent)) is not None: output_type_name = output_schema.name() else: output_type_name = "str" current_span = agent_span( name=current_agent.name, handoffs=[], tools=[], output_type=output_type_name, ) current_span.start(mark_as_current=True) current_turn += 1 if max_turns is not None and current_turn > max_turns: _error_tracing.attach_error_to_span( current_span, SpanError( message="Max turns exceeded", data={"max_turns": max_turns}, ), ) max_turns_error = MaxTurnsExceeded(f"Max turns ({max_turns}) exceeded") run_error_data = build_run_error_data( input=original_input, new_items=session_items, raw_responses=model_responses, last_agent=current_agent, reasoning_item_id_policy=resolved_reasoning_item_id_policy, ) handler_result = await resolve_run_error_handler_result( error_handlers=error_handlers, error_kind="max_turns", error=max_turns_error, context_wrapper=context_wrapper, run_data=run_error_data, ) if handler_result is None: raise max_turns_error include_in_history = handler_result.include_in_history handler_output_recorded = False handler_persisted_item_count = 0 async def _save_max_turns_handler_output( items: list[RunItem], store_setting: bool | None = store_setting, generated_items: list[RunItem] = generated_items, session_items: list[RunItem] = session_items, ) -> None: nonlocal handler_output_recorded, handler_persisted_item_count handler_persisted_item_count = ( await save_final_turn_items_after_guardrails( session=session, run_state=None, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=_attempt_input_guardrail_results(), items=items, response_id=None, reasoning_item_id_policy=resolved_reasoning_item_id_policy, store=store_setting, wrapper=context_wrapper, ) ) if not items: return generated_items.extend(items) session_items.extend(items) handler_output_recorded = True ( validated_output, synthesized_item, ) = await finalize_max_turns_handler_output( agent=current_agent, hooks=hooks, run_config=run_config, output=handler_result.final_output, context_wrapper=context_wrapper, output_guardrail_results=output_guardrail_results, save_items_after_guardrails=_save_max_turns_handler_output, include_in_history=include_in_history, ) if include_in_history and not handler_output_recorded: await _save_max_turns_handler_output([synthesized_item]) current_step = getattr(run_state, "_current_step", None) approvals_from_state = approvals_from_step(current_step) result = RunResult( input=original_input, new_items=session_items, raw_responses=model_responses, final_output=validated_output, _last_agent=current_agent, input_guardrail_results=input_guardrail_results, output_guardrail_results=output_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, context_wrapper=context_wrapper, interruptions=approvals_from_state, _tool_use_tracker_snapshot=_tool_use_tracker_snapshot(), max_turns=max_turns, ) result._current_turn = max_turns result._model_input_items = list(generated_items) result._replay_from_model_input_items = list(generated_items) != list( session_items ) if run_state is not None: result._trace_state = run_state._trace_state result._current_turn_persisted_item_count = handler_persisted_item_count result._original_input = copy_input_items(original_input) return _finalize_result(result) if run_state is not None and ( not resuming_turn or isinstance(run_state._current_step, NextStepRunAgain) ): run_state._current_turn_persisted_item_count = 0 logger.debug("Running agent %s (turn %s)", current_agent.name, current_turn) if session_persistence_enabled: try: last_saved_input_snapshot_for_rewind = ( ItemHelpers.input_to_new_input_list(original_input) ) except Exception: last_saved_input_snapshot_for_rewind = None if run_state is not None and _has_output_guardrails(current_agent, run_config): _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, ) blocked_output_owner_starts = _BlockedOutputOwnerStarts( nonstreamed_session_items=len(session_items), run_state_generated_items=( len(run_state._generated_items) if run_state is not None else None ), run_state_session_items=( len(run_state._session_items) if run_state is not None else None ), run_state_model_responses=( len(run_state._model_responses) if run_state is not None else None ), run_state_tool_output_guardrail_results=( len(run_state._tool_output_guardrail_results) if run_state is not None else None ), ) items_for_model = ( pending_server_items if server_conversation_tracker is not None and pending_server_items else generated_items ) turn_usage_start = snapshot_usage(context_wrapper.usage) current_turn_span = ( turn_span( turn=current_turn, agent_name=current_agent.name, ) if use_task_and_turn_spans else None ) if current_turn_span is not None: current_turn_span.start(mark_as_current=True) try: if current_turn <= 1: try: if sequential_guardrails: await run_input_guardrails( starting_agent, sequential_guardrails, copy_input_items(original_input), context_wrapper, input_guardrail_results, ) except InputGuardrailTripwireTriggered: session_input_items_for_persistence = ( await persist_session_items_for_guardrail_trip( session, server_conversation_tracker, session_input_items_for_persistence, original_user_input, run_state, store=store_setting, wrapper=context_wrapper, ) ) raise model_task = asyncio.create_task( run_single_turn( bindings=current_bindings, all_tools=all_tools, 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, ) ) if parallel_guardrails: guardrail_task = asyncio.create_task( run_input_guardrails( starting_agent, parallel_guardrails, copy_input_items(original_input), context_wrapper, input_guardrail_results, ) ) try: _, turn_result = await asyncio.gather( guardrail_task, model_task, ) except InputGuardrailTripwireTriggered: if should_cancel_parallel_model_task_on_input_guardrail_trip(): if not model_task.done(): model_task.cancel() await asyncio.gather(model_task, return_exceptions=True) session_input_items_for_persistence = ( await persist_session_items_for_guardrail_trip( session, server_conversation_tracker, session_input_items_for_persistence, original_user_input, run_state, store=store_setting, wrapper=context_wrapper, ) ) 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 ) raise else: turn_result = await model_task else: turn_result = await run_single_turn( bindings=current_bindings, all_tools=all_tools, 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, ) finally: if current_turn_span is not None: attach_usage_to_span( current_turn_span, usage_delta(turn_usage_start, context_wrapper.usage), ) current_turn_span.finish(reset_current=True) # Start hooks should only run on the first turn unless reset by a handoff. last_saved_input_snapshot_for_rewind = None should_run_agent_start_hooks = False model_responses.append(turn_result.model_response) input_before_turn_rewrite = original_input original_input = turn_result.original_input # For model input, use new_step_items (filtered on handoffs). generated_items = turn_result.pre_step_items + turn_result.new_step_items # Accumulate unfiltered items for observability. turn_session_items = session_items_for_turn(turn_result) session_items.extend(turn_session_items) if pending_input_admission_items and run_state is not None: run_state._generated_items = list(generated_items) run_state._session_items = list(session_items) run_state._model_responses = list(model_responses) run_state._last_processed_response = turn_result.processed_response run_state._current_turn = current_turn run_state._mark_generated_items_merged_with_last_processed() pending_input_admission_items = [] if run_state is not None and turn_result.nested_history_owned_items is not None: run_state._nested_history_owned_session_item_refs = ( reconcile_nested_history_owned_session_item_refs( session_items, run_state._nested_history_owned_session_item_refs, input_before_turn_rewrite, turn_result.original_input, turn_result.nested_history_owned_items, ) ) if server_conversation_tracker is not None: pending_server_items = list(turn_result.new_step_items) server_conversation_tracker.track_server_items(turn_result.model_response) tool_input_guardrail_results.extend(turn_result.tool_input_guardrail_results) tool_output_guardrail_results.extend(turn_result.tool_output_guardrail_results) items_to_save_turn = list(turn_session_items) if not isinstance(turn_result.next_step, NextStepInterruption): if session_persistence_enabled: output_call_ids = { item.raw_item.get("call_id") if isinstance(item.raw_item, dict) else getattr(item.raw_item, "call_id", None) for item in turn_result.new_step_items if item.type == "tool_call_output_item" and ( item.raw_item.get("type") if isinstance(item.raw_item, dict) else getattr(item.raw_item, "type", None) ) != "program_output" } for item in generated_items: if item.type != "tool_call_item": continue call_id = ( item.raw_item.get("call_id") if isinstance(item.raw_item, dict) else getattr(item.raw_item, "call_id", None) ) if ( call_id in output_call_ids and item not in items_to_save_turn and not ( run_state is not None and run_state._current_turn_persisted_item_count > 0 ) ): items_to_save_turn.append(item) if items_to_save_turn and not isinstance( turn_result.next_step, NextStepFinalOutput ): logger.debug( "Persisting turn items (types=%s)", [item.type for item in items_to_save_turn], ) if is_resumed_state and run_state is not None: saved_count = await save_result_to_session( session, [], items_to_save_turn, None, response_id=turn_result.model_response.response_id, reasoning_item_id_policy=( run_state._reasoning_item_id_policy ), store=store_setting, wrapper=context_wrapper, ) run_state._current_turn_persisted_item_count += saved_count else: await save_result_to_session( session, [], items_to_save_turn, run_state, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) # After the first resumed turn, treat subsequent turns as fresh # so counters and input saving behave normally. is_resumed_state = False try: if isinstance(turn_result.next_step, NextStepFinalOutput): if run_state is not None and _has_output_guardrails( current_agent, run_config ): run_state._tool_output_guardrail_results = list( tool_output_guardrail_results ) output_guardrail_result_start = len(output_guardrail_results) try: await run_output_guardrails( current_agent.output_guardrails + (run_config.output_guardrails or []), current_agent, turn_result.next_step.output, context_wrapper, output_guardrail_results, ) except OutputGuardrailTripwireTriggered as exc: if not _is_terminal_tool_output_response( turn_session_items, turn_result.processed_response, run_state, ): raise sanitized_results = _sanitize_blocked_output_guardrail_results( output_guardrail_results[output_guardrail_result_start:], exc ) output_guardrail_results[output_guardrail_result_start:] = ( sanitized_results ) session_items = _blocked_output_failure_items( session_items, (), blocked_output_owner_starts, ) blocked_message = _resolve_output_guardrail_blocked_message( exc, agent=current_agent, run_config=run_config, context_wrapper=context_wrapper, ) if blocked_message != OUTPUT_GUARDRAIL_BLOCKED_TOOL_OUTPUT: sanitized_results = _sanitize_blocked_output_guardrail_results( sanitized_results, exc, blocked_message, ) output_guardrail_results[output_guardrail_result_start:] = ( sanitized_results ) retained_items = _retained_items_for_blocked_response( turn_session_items, turn_result.model_response, run_state, turn_result.processed_response, owner_starts=blocked_output_owner_starts, blocked_message=blocked_message, ) list.extend(session_items, retained_items) try: await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=(session_persistence_enabled), input_guardrail_results=( _attempt_input_guardrail_results() ), items=retained_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) except BaseException as persistence_error: raise _safe_redacted_persistence_error( persistence_error ) from None raise except (Exception, asyncio.CancelledError) as guardrail_error: if not isinstance( guardrail_error, asyncio.CancelledError ) or not _is_terminal_tool_output_response( turn_session_items, turn_result.processed_response, run_state, ): final_turn_items = _final_turn_items_for_persistence( turn_session_items, turn_result.processed_response, run_state, current_agent, run_config, ) await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=( _attempt_input_guardrail_results() ), items=final_turn_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) raise final_turn_items = _final_turn_items_for_persistence( turn_session_items, turn_result.processed_response, run_state, current_agent, run_config, ) await save_final_turn_items_after_guardrails( session=session, run_state=run_state, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=_attempt_input_guardrail_results(), items=final_turn_items, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) # Ensure starting_input is not None and not RunState final_output_result_input: str | list[TResponseInputItem] = ( normalized_starting_input ) result = RunResult( input=final_output_result_input, new_items=session_items, raw_responses=model_responses, final_output=turn_result.next_step.output, _last_agent=current_agent, input_guardrail_results=input_guardrail_results, output_guardrail_results=output_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, context_wrapper=context_wrapper, interruptions=[], _tool_use_tracker_snapshot=_tool_use_tracker_snapshot(), max_turns=max_turns, ) result._current_turn = current_turn result._model_input_items = list(generated_items) result._replay_from_model_input_items = list(generated_items) != list( session_items ) if run_state is not None: result._current_turn_persisted_item_count = ( run_state._current_turn_persisted_item_count ) result._original_input = copy_input_items(original_input) if run_state is not None: run_state._current_step = None return _finalize_result(result) elif isinstance(turn_result.next_step, NextStepInterruption): if session_persistence_enabled and not ( _should_defer_interrupted_session_items( current_agent, run_config, ) ): if not input_guardrails_triggered( _attempt_input_guardrail_results() ): # Persist session items but skip approval placeholders. input_items_for_save_interruption: list[TResponseInputItem] = ( session_input_items_for_persistence if session_input_items_for_persistence is not None else [] ) await save_result_to_session( session, input_items_for_save_interruption, session_items_for_turn(turn_result), run_state, response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) append_model_response_if_new( model_responses, turn_result.model_response ) processed_response_for_state = resolve_processed_response( run_state=run_state, processed_response=turn_result.processed_response, ) if run_state is not None: update_run_state_for_interruption( run_state=run_state, model_responses=model_responses, processed_response=processed_response_for_state, generated_items=generated_items, session_items=session_items, current_turn=current_turn, next_step=turn_result.next_step, ) # Ensure starting_input is not None and not RunState interruption_result_input2: str | list[TResponseInputItem] = ( normalized_starting_input ) result = build_interruption_result( result_input=interruption_result_input2, session_items=session_items, model_responses=model_responses, current_agent=current_agent, input_guardrail_results=input_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, context_wrapper=context_wrapper, interruptions=approvals_from_step(turn_result.next_step), processed_response=processed_response_for_state, tool_use_tracker=tool_use_tracker, max_turns=max_turns, current_turn=current_turn, generated_items=generated_items, run_state=run_state, original_input=original_input, ) return _finalize_result(result) elif isinstance(turn_result.next_step, NextStepHandoff): current_agent = cast(Agent[TContext], turn_result.next_step.new_agent) if run_state is not None: run_state._current_agent = current_agent # Next agent starts with the nested/filtered input. # Assign without type annotation to avoid redefinition error starting_input = turn_result.original_input original_input = turn_result.original_input current_span.finish(reset_current=True) current_span = None should_run_agent_start_hooks = True elif isinstance(turn_result.next_step, NextStepRunAgain): await save_turn_items_if_needed( session=session, run_state=run_state, session_persistence_enabled=session_persistence_enabled, input_guardrail_results=input_guardrail_results, items=session_items_for_turn(turn_result), response_id=turn_result.model_response.response_id, store=store_setting, wrapper=context_wrapper, ) continue else: raise AgentsException( f"Unknown next step type: {type(turn_result.next_step)}" ) finally: # execute_tools_and_side_effects returns a SingleStepResult that # stores direct references to the `pre_step_items` and `new_step_items` # lists it manages internally. Clear them here so the next turn does not # hold on to items from previous turns and to avoid leaking agent refs. turn_result.pre_step_items.clear() turn_result.new_step_items.clear() except BaseException as exc: run_exception = exc if _is_error_data_redacted(exc): _detach_data_redacted_error_traceback(exc) else: attach_generic_agent_error( current_span, exc, trace_include_sensitive_data=run_config.trace_include_sensitive_data, ) if isinstance(exc, AgentsException): _clear_data_redacted_error_traceback(exc) exc.run_data = RunErrorDetails( input=original_input, new_items=session_items, raw_responses=model_responses, last_agent=current_agent, context_wrapper=context_wrapper, input_guardrail_results=input_guardrail_results, output_guardrail_results=output_guardrail_results, tool_input_guardrail_results=tool_input_guardrail_results, tool_output_guardrail_results=tool_output_guardrail_results, ) raise finally: await cleanup_models_after_run(tool_use_tracker) try: try: memory_input = _sandbox_memory_input( memory_input_items_for_persistence=memory_input_items_for_persistence, original_user_input=original_user_input, original_input=original_input, ) if completed_result is not None: await sandbox_runtime.enqueue_memory_result( completed_result, input_override=memory_input, ) elif run_exception is not None: current_step = getattr(run_state, "_current_step", None) await sandbox_runtime.enqueue_memory_payload( input=memory_input, new_items=session_items, final_output=None, interruptions=approvals_from_step(current_step), terminal_metadata=terminal_metadata_for_exception(run_exception), ) except Exception as error: log_model_and_tool_action_warning( logger, "Failed to enqueue sandbox memory after run", error ) sandbox_resume_state = await sandbox_runtime.cleanup() except Exception as error: log_tool_action_warning( logger, "Failed to clean up sandbox resources after run", error ) else: if completed_result is not None: completed_result._sandbox_resume_state = sandbox_resume_state finally: if completed_result is not None: completed_result._sandbox_session = None try: await dispose_resolved_computers(run_context=context_wrapper) except Exception as error: log_tool_action_warning(logger, "Failed to dispose computers after run", error) if current_span is not None: current_span.finish(reset_current=True) if current_task_span is not None: attach_usage_to_span( current_task_span, usage_delta(task_usage_start, context_wrapper.usage), ) current_task_span.finish(reset_current=True) def run_sync( self, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], **kwargs: Unpack[RunOptions[TContext]], ) -> RunResult: redacted_error: BaseException | None = None redacted_source: BaseException | None try: return self._run_sync_impl(starting_agent, input, **kwargs) except BaseException as error: if _is_error_data_redacted(error): redacted_source = error else: redacted_source = _data_redacted_sync_cancellation_source(error) if redacted_source is None: raise if isinstance(redacted_source, asyncio.CancelledError): redacted_error = _prepare_data_redacted_error(redacted_source) else: _detach_data_redacted_error_traceback(redacted_source) redacted_error = redacted_source self = cast(Any, None) starting_agent = cast(Any, None) input = cast(Any, None) cast(dict[str, Any], kwargs).clear() assert redacted_error is not None _detach_data_redacted_error_traceback(redacted_error) _raise_data_redacted_error(redacted_error) def _run_sync_impl( self, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], **kwargs: Unpack[RunOptions[TContext]], ) -> RunResult: context = kwargs.get("context") max_turns = kwargs.get("max_turns", DEFAULT_MAX_TURNS) hooks = kwargs.get("hooks") run_config = kwargs.get("run_config") error_handlers = kwargs.get("error_handlers") previous_response_id = kwargs.get("previous_response_id") auto_previous_response_id = kwargs.get("auto_previous_response_id", False) conversation_id = kwargs.get("conversation_id") session = kwargs.get("session") # Python 3.14 stopped implicitly wiring up a default event loop # when synchronous code touches asyncio APIs for the first time. # Several of our synchronous entry points (for example the Redis/SQLAlchemy session helpers) # construct asyncio primitives like asyncio.Lock during __init__, # which binds them to whatever loop happens to be the thread's default at that moment. # To keep those locks usable we must ensure that run_sync reuses that same default loop # instead of hopping over to a brand-new asyncio.run() loop. try: already_running_loop = asyncio.get_running_loop() except RuntimeError: already_running_loop = None if already_running_loop is not None: # This method is only expected to run when no loop is already active. # (Each thread has its own default loop; concurrent sync runs should happen on # different threads. In a single thread use the async API to interleave work.) raise RuntimeError( "AgentRunner.run_sync() cannot be called when an event loop is already running." ) policy = asyncio.get_event_loop_policy() with warnings.catch_warnings(): warnings.simplefilter("ignore", DeprecationWarning) try: default_loop = policy.get_event_loop() except RuntimeError: default_loop = policy.new_event_loop() policy.set_event_loop(default_loop) if default_loop.is_closed(): default_loop = policy.new_event_loop() policy.set_event_loop(default_loop) # We intentionally leave the default loop open even if we had to create one above. Session # instances and other helpers stash loop-bound primitives between calls and expect to find # the same default loop every time run_sync is invoked on this thread. # Schedule the async run on the default loop so that we can manage cancellation explicitly. task = default_loop.create_task( self.run( starting_agent, input, session=session, context=context, max_turns=max_turns, hooks=hooks, run_config=run_config, error_handlers=error_handlers, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, conversation_id=conversation_id, ) ) try: # Drive the coroutine to completion, harvesting the final RunResult. return default_loop.run_until_complete(task) except BaseException as error: # If the sync caller aborts (KeyboardInterrupt, etc.), make sure the scheduled task # does not linger on the shared loop by cancelling it and waiting for completion. if not task.done(): task.cancel() with contextlib.suppress(asyncio.CancelledError): default_loop.run_until_complete(task) if _is_error_data_redacted(error) or isinstance(error, ModelBehaviorError): _detach_data_redacted_error_traceback(error) raise finally: if not default_loop.is_closed(): # The loop stays open for subsequent runs, but we still need to flush any pending # async generators so their cleanup code executes promptly. with contextlib.suppress(RuntimeError): default_loop.run_until_complete(default_loop.shutdown_asyncgens()) def run_streamed( self, starting_agent: Agent[TContext], input: str | list[TResponseInputItem] | RunState[TContext], **kwargs: Unpack[RunOptions[TContext]], ) -> RunResultStreaming: context = kwargs.get("context") max_turns = kwargs.get("max_turns", DEFAULT_MAX_TURNS) hooks = cast(RunHooks[TContext], validate_run_hooks(kwargs.get("hooks"))) run_config = kwargs.get("run_config") error_handlers = kwargs.get("error_handlers") previous_response_id = kwargs.get("previous_response_id") auto_previous_response_id = kwargs.get("auto_previous_response_id", False) conversation_id = kwargs.get("conversation_id") session = kwargs.get("session") run_config = RunConfig() if run_config is None else _coerce_run_config(run_config) # Handle RunState input is_resumed_state = isinstance(input, RunState) run_state: RunState[TContext] | None = None input_for_result: str | list[TResponseInputItem] starting_input = input if not is_resumed_state else None if is_resumed_state: run_state = cast(RunState[TContext], input) ( conversation_id, previous_response_id, auto_previous_response_id, ) = apply_resumed_conversation_settings( run_state=run_state, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) validate_session_conversation_settings( session, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) # When resuming, use the original_input from state. # primeFromState will mark items as sent so prepareInput skips them starting_input = run_state._original_input logger.debug( "Resuming from RunState in run_streaming()", extra=build_resumed_stream_debug_extra( run_state, include_tool_output=not _debug.DONT_LOG_TOOL_DATA, ), ) # When resuming, use the original_input from state. # primeFromState will mark items as sent so prepareInput skips them raw_input_for_result = run_state._original_input input_for_result = normalize_resumed_input(raw_input_for_result) ( input_for_result, run_state._nested_history_owned_session_item_refs, ) = reconcile_nested_history_owned_input_after_rewrite( raw_input_for_result, input_for_result, run_state._nested_history_owned_session_item_refs, ) run_state._original_input = copy_input_items(input_for_result) # Use context from RunState if not provided, otherwise override it. context_wrapper = resolve_resumed_context( run_state=run_state, context=context, ) context = context_wrapper.context # Override max_turns with the state's max_turns to preserve it across resumption max_turns = run_state._max_turns else: # input is already str | list[TResponseInputItem] when not RunState # Reuse input_for_result variable from outer scope input_for_result = cast(str | list[TResponseInputItem], input) validate_session_conversation_settings( session, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) context_wrapper = ensure_context_wrapper(context) set_agent_tool_state_scope(context_wrapper, None) # input_for_state is the same as input_for_result here input_for_state = input_for_result run_state = RunState( context=context_wrapper, original_input=copy_input_items(input_for_state), starting_agent=starting_agent, max_turns=max_turns, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) resolved_reasoning_item_id_policy: ReasoningItemIdPolicy | None = ( run_config.reasoning_item_id_policy if run_config.reasoning_item_id_policy is not None else (run_state._reasoning_item_id_policy if run_state is not None else None) ) if run_state is not None: run_state._reasoning_item_id_policy = resolved_reasoning_item_id_policy schema_agent = ( run_state._current_agent if run_state is not None and run_state._current_agent is not None else starting_agent ) validate_output_guardrails_with_server_managed_conversation( schema_agent, run_config, conversation_id=conversation_id, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, ) ( trace_workflow_name, trace_id, trace_group_id, trace_metadata, trace_config, ) = resolve_trace_settings(run_state=run_state, run_config=run_config) # If there's already a trace, we don't create a new one. In addition, we can't end the # trace here, because the actual work is done in `stream_events` and this method ends # before that. new_trace = create_trace_for_run( workflow_name=trace_workflow_name, trace_id=trace_id, group_id=trace_group_id, metadata=trace_metadata, tracing=trace_config, disabled=run_config.tracing_disabled, trace_state=run_state._trace_state if run_state is not None else None, reattach_resumed_trace=is_resumed_state, ) if run_state is not None: run_state.set_trace(new_trace if new_trace is not None else get_current_trace()) sandbox_runtime = SandboxRuntime( starting_agent=starting_agent, run_config=run_config, rollout_id=_sandbox_memory_rollout_id( run_config=run_config, conversation_id=conversation_id, session=session, ), run_state=run_state, ) sandbox_runtime.assert_agent_supported(schema_agent) output_schema = get_output_schema(schema_agent) streamed_input: str | list[TResponseInputItem] = ( starting_input if starting_input is not None and not isinstance(starting_input, RunState) else "" ) streamed_result = RunResultStreaming( input=copy_input_items(streamed_input), # When resuming from RunState, use session_items from state. # primeFromState will mark items as sent so prepareInput skips them. # Copy it: the streamed loop appends to new_items, and the caller still # owns the state as a resumable snapshot. new_items=list(run_state._session_items) if run_state is not None else [], current_agent=schema_agent, raw_responses=run_state._model_responses if run_state is not None else [], final_output=None, is_complete=False, current_turn=run_state._current_turn if run_state is not None else 0, max_turns=max_turns, input_guardrail_results=( list(run_state._input_guardrail_results) if run_state is not None else [] ), output_guardrail_results=( list(run_state._output_guardrail_results) if run_state is not None else [] ), tool_input_guardrail_results=( list(getattr(run_state, "_tool_input_guardrail_results", [])) if run_state is not None else [] ), tool_output_guardrail_results=( list(getattr(run_state, "_tool_output_guardrail_results", [])) if run_state is not None else [] ), _current_agent_output_schema=output_schema, trace=new_trace, context_wrapper=context_wrapper, interruptions=[], # Preserve persisted-count from state to avoid re-saving items when resuming. # If a cross-SDK state omits the counter, fall back to len(generated_items) # to avoid duplication. _current_turn_persisted_item_count=( run_state._current_turn_persisted_item_count if run_state is not None else 0 ), # When resuming from RunState, preserve the original input from the state # This ensures originalInput in serialized state reflects the first turn's input _original_input=( copy_input_items(run_state._original_input) if run_state is not None and run_state._original_input is not None else copy_input_items(streamed_input) ), ) streamed_result._model_input_items = ( list(run_state._generated_items) if run_state is not None else [] ) streamed_result._replay_from_model_input_items = ( list(run_state._generated_items) != list(run_state._session_items) if run_state is not None else False ) streamed_result._reasoning_item_id_policy = resolved_reasoning_item_id_policy if run_state is not None: streamed_result._trace_state = run_state._trace_state # Store run_state in streamed_result._state so it's accessible throughout streaming # Now that we create run_state for both fresh and resumed runs, always set it streamed_result._conversation_id = conversation_id streamed_result._previous_response_id = previous_response_id streamed_result._auto_previous_response_id = auto_previous_response_id streamed_result._state = run_state if run_state is not None: streamed_result._tool_use_tracker_snapshot = run_state.get_tool_use_tracker_snapshot() if sandbox_runtime.enabled: sandbox_runtime.apply_result_metadata(streamed_result) # Kick off the actual agent loop in the background and return the streamed result object. streamed_result.run_loop_task = asyncio.create_task( _await_data_redacted_error_boundary( lambda: start_streaming( starting_input=input_for_result, streamed_result=streamed_result, starting_agent=starting_agent, max_turns=max_turns, hooks=hooks, context_wrapper=context_wrapper, run_config=run_config, error_handlers=error_handlers, previous_response_id=previous_response_id, auto_previous_response_id=auto_previous_response_id, conversation_id=conversation_id, session=session, run_state=run_state, trace_workflow_name=trace_workflow_name, is_resumed_state=is_resumed_state, sandbox_runtime=sandbox_runtime, ) ) ) if sandbox_runtime.enabled: streamed_result.ensure_sandbox_cleanup_on_completion() return streamed_result DEFAULT_AGENT_RUNNER = AgentRunner()