Skip to content
Merged
7 changes: 6 additions & 1 deletion src/agents/memory/openai_responses_compaction_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -431,7 +431,6 @@ async def _run_compaction_locked(
if is_automatic and snapshot is None and self.max_rollback_items is not None:
await self._get_all_underlying_session_items()

self._deferred_response_id = None
logger.debug(
"compact: start for %s using %s (mode=%s)",
self._response_id,
Expand Down Expand Up @@ -490,6 +489,12 @@ async def _run_compaction_locked(
)
self._session_items = None if read_items is not None else output_items

# Clear the deferred marker only now that compaction has actually settled. Clearing it
# before the fallible API call/replacement above would let a failed forced compaction
# silently lose its "this must be forced" signal: a later retry recomputes `force` from
# this marker, so an early clear makes the retry decline work that was still owed.
self._deferred_response_id = None

logger.debug(
"compact: done for %s (mode=%s, output=%s, candidates=%s)",
self._response_id,
Expand Down
14 changes: 6 additions & 8 deletions src/agents/run_internal/agent_runner_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -505,18 +505,16 @@ def build_interruption_result(


def reject_unrecoverable_terminal_state(run_state: RunState | None) -> None:
"""Fail closed when a previous run already produced a final output that cannot be reproduced.
"""Fail closed when a previous run reached a boundary that cannot safely be resumed.

The marker is set once that output, its guardrails, and its terminal hooks have completed,
and is cleared only once the turn is fully persisted. In between, the run owns a result no
resume can settle, so resuming would repeat the model call and the lifecycle hooks for an
output the caller already received. Raised before any Session, sandbox, model, tool,
guardrail, or hook work so the rejection has no side effects of its own.
This includes failed terminal persistence and a speculative handoff rejected by an input
guardrail. Resuming would repeat completed work or bypass the failed input check. Reject
before any Session, sandbox, model, tool, guardrail, or hook work.
"""
if run_state is not None and run_state._terminal_unrecoverable:
raise UserError(
"This RunState already produced a final output whose Session write did not "
"complete, so it cannot be resumed. Start a new run instead."
"This RunState ended at an unrecoverable boundary and cannot be resumed. "
"Start a new run instead."
)


Expand Down
48 changes: 38 additions & 10 deletions src/agents/run_internal/run_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
_detach_data_redacted_error_traceback,
_is_error_data_redacted,
_mark_error_data_redacted,
_mark_error_to_drain_stream_events,
_prepare_data_redacted_error,
)
from ..guardrail import OutputGuardrailResult
Expand Down Expand Up @@ -422,6 +423,7 @@ async def _save_stream_items(
response_id: str | None,
update_persisted_count: bool,
store: bool | None = None,
resumed_write_state: RunState | None = None,
) -> None:
if not await _should_persist_stream_items(
session=session,
Expand All @@ -437,6 +439,7 @@ async def _save_stream_items(
response_id=response_id,
store=store,
wrapper=streamed_result.context_wrapper,
resumed_write_state=resumed_write_state,
)
if update_persisted_count and streamed_result._state is not None:
streamed_result._current_turn_persisted_item_count = (
Expand Down Expand Up @@ -1165,6 +1168,12 @@ async def _save_stream_items_without_count(
response_id=response_id,
update_persisted_count=False,
store=store_setting,
resumed_write_state=(
run_state
if run_state is not None
and isinstance(run_state._current_step, NextStepRunAgain)
else None
Comment thread
seratch marked this conversation as resolved.
),
)

async def _save_max_turns_items(
Expand Down Expand Up @@ -1886,23 +1895,42 @@ def _record_max_turns_handler_output(
server_conversation_tracker.track_server_items(turn_result.model_response)

if isinstance(turn_result.next_step, NextStepHandoff):
await _save_stream_items_without_count(
turn_session_items,
turn_result.model_response.response_id,
store_setting,
)
# A failed input check makes this completed speculative turn non-resumable:
# replaying it would skip the starting agent's guardrails. Keep its executed
# call/output records coherent instead of partially rolling back history.
try:
if streamed_result._input_guardrails_task is not None:
await streamed_result._input_guardrails_task
for guardrail_result in streamed_result.input_guardrail_results:
if guardrail_result.output.tripwire_triggered:
raise InputGuardrailTripwireTriggered(guardrail_result)
except BaseException:
if run_state is not None:
run_state._terminal_unrecoverable = True
raise
current_agent = turn_result.next_step.new_agent
if run_state is not None:
run_state._current_agent = current_agent
_publish_streamed_result_agent(streamed_result, current_agent)
Comment thread
seratch marked this conversation as resolved.
Comment thread
seratch marked this conversation as resolved.
current_span.finish(reset_current=True)
current_span = None
should_run_agent_start_hooks = True
if streamed_result._state is not None:
streamed_result._state._current_step = NextStepRunAgain()
Comment thread
seratch marked this conversation as resolved.
# Queue the agent-transition event before the fallible session append so
# stream consumers observe the transition even if the append later raises.
streamed_result._event_queue.put_nowait(
AgentUpdatedStreamEvent(new_agent=current_agent)
)
Comment thread
seratch marked this conversation as resolved.
if streamed_result._state is not None:
streamed_result._state._current_step = NextStepRunAgain()
try:
await _save_stream_items_without_count(
turn_session_items,
turn_result.model_response.response_id,
store_setting,
)
except BaseException as session_persistence_error:
_mark_error_to_drain_stream_events(session_persistence_error)
raise
current_span.finish(reset_current=True)
current_span = None
should_run_agent_start_hooks = True

if await _wait_for_streamed_turn_events_and_stop_if_cancelled(streamed_result):
break
Expand Down
150 changes: 100 additions & 50 deletions src/agents/run_internal/session_persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -631,6 +631,71 @@ def update_run_state_after_resume(
run_state._current_step = next_step # type: ignore[assignment]


async def _apply_post_write_compaction(
session: Session,
*,
response_id: str | None,
store: bool | None,
has_local_tool_outputs: bool,
wrapper: RunContextWrapper[Any] | None = None,
) -> None:
"""Evaluate deferred/forced Responses compaction for a settled session append.

Shared by the immediate-write path in ``save_result_to_session`` and the checkpoint
replay path in ``resume_pending_session_write``, so a batch that only settles later
(via a separate resume) still gets the same compaction decision it would have gotten
had the original append succeeded inline. ``wrapper`` is the caller's raw (pre-gating)
context wrapper; it is used as-is for ``run_compaction`` and re-gated here for
``_defer_compaction``, mirroring the two call sites this helper replaces.
"""
if not response_id or not is_openai_responses_compaction_aware_session(session):
return

if has_local_tool_outputs:
defer_compaction = getattr(session, "_defer_compaction", None)
if callable(defer_compaction):
await _call_session_method(
defer_compaction,
response_id,
store=store,
wrapper=_get_session_wrapper(session, wrapper),
)
logger.debug(
"skip: deferring compaction for response %s due to local tool outputs",
response_id,
)
return

deferred_response_id = None
get_deferred = getattr(session, "_get_deferred_compaction_response_id", None)
if callable(get_deferred):
deferred_response_id = get_deferred()
force_compaction = deferred_response_id is not None
if force_compaction:
logger.debug(
"compact: forcing for response %s after deferred %s",
response_id,
deferred_response_id,
)
compaction_args: OpenAIResponsesCompactionArgs = {
"response_id": response_id,
"force": force_compaction,
}
if store is not None:
compaction_args["store"] = store
if wrapper is not None:
wrapper._session_compaction_is_automatic = True # type: ignore[attr-defined]
try:
await _call_session_method(
session.run_compaction,
compaction_args,
wrapper=wrapper,
)
finally:
if wrapper is not None:
wrapper._session_compaction_is_automatic = False # type: ignore[attr-defined]


async def save_result_to_session(
session: Session | None,
original_input: str | list[TResponseInputItem],
Expand Down Expand Up @@ -738,6 +803,10 @@ async def save_result_to_session(
run_state._current_turn_persisted_item_count = already_persisted + saved_run_items_count
return saved_run_items_count

has_local_tool_outputs = any(
isinstance(item, ToolCallOutputItem | HandoffOutputItem) for item in new_items
)

if resumed_write_state is not None:
if resumed_write_state._pending_session_write is not None:
raise UserError("Resolve the pending Session write before saving another batch")
Expand All @@ -748,7 +817,13 @@ async def save_result_to_session(
"persisted_count": (
resumed_write_state._current_turn_persisted_item_count + saved_run_items_count
),
"response_id": response_id,
"store": store,
"has_local_tool_outputs": has_local_tool_outputs,
}
# resume_pending_session_write() applies post-write compaction itself once the
# checkpoint settles, whether that happens inline below or on a later, separate
# resume -- so it is not repeated after this call returns.
await resume_pending_session_write(
resumed_write_state,
session,
Expand All @@ -760,53 +835,14 @@ async def save_result_to_session(
if run_state is not None:
run_state._current_turn_persisted_item_count = already_persisted + saved_run_items_count

if response_id and is_openai_responses_compaction_aware_session(session):
has_local_tool_outputs = any(
isinstance(item, ToolCallOutputItem | HandoffOutputItem) for item in new_items
if resumed_write_state is None:
await _apply_post_write_compaction(
session,
response_id=response_id,
store=store,
has_local_tool_outputs=has_local_tool_outputs,
wrapper=compaction_wrapper,
)
if has_local_tool_outputs:
defer_compaction = getattr(session, "_defer_compaction", None)
if callable(defer_compaction):
await _call_session_method(
defer_compaction,
response_id,
store=store,
wrapper=wrapper,
)
logger.debug(
"skip: deferring compaction for response %s due to local tool outputs",
response_id,
)
return saved_run_items_count

deferred_response_id = None
get_deferred = getattr(session, "_get_deferred_compaction_response_id", None)
if callable(get_deferred):
deferred_response_id = get_deferred()
force_compaction = deferred_response_id is not None
if force_compaction:
logger.debug(
"compact: forcing for response %s after deferred %s",
response_id,
deferred_response_id,
)
compaction_args: OpenAIResponsesCompactionArgs = {
"response_id": response_id,
"force": force_compaction,
}
if store is not None:
compaction_args["store"] = store
if compaction_wrapper is not None:
compaction_wrapper._session_compaction_is_automatic = True # type: ignore[attr-defined]
try:
await _call_session_method(
session.run_compaction,
compaction_args,
wrapper=compaction_wrapper,
)
finally:
if compaction_wrapper is not None:
compaction_wrapper._session_compaction_is_automatic = False # type: ignore[attr-defined]

return saved_run_items_count

Expand Down Expand Up @@ -886,10 +922,10 @@ def digests(items: Sequence[TResponseInputItem]) -> list[str]:
append = True
else:
expected = before + digests(pending["items"])
committed_generation: int | None = None
observed_generation: int | None = None
get_with_generation = getattr(session, "_get_items_with_generation", None)
if wrapper is not None and callable(get_with_generation):
tail, committed_generation = await _call_session_method(
tail, observed_generation = await _call_session_method(
get_with_generation,
lambda: _session_get_items(session, limit=len(expected), wrapper=wrapper),
)
Expand All @@ -904,12 +940,26 @@ def digests(items: Sequence[TResponseInputItem]) -> list[str]:
"Repair the original Session before resuming; do not rerun the completed tool."
)
append = unchanged
if committed and committed_generation is not None and wrapper is not None:
wrapper._session_compaction_generation = committed_generation # type: ignore[attr-defined]
# The original append can advance the wrapper generation even when it fails
# atomically. Reconciled unchanged history is also safe to append against;
# subsequent mutations still revoke ownership through the normal generation check.
if observed_generation is not None and wrapper is not None:
wrapper._session_compaction_generation = observed_generation # type: ignore[attr-defined]
if append:
# Backends may retain or transform their input; the durable checkpoint stays detached.
await _session_add_items(session, copy.deepcopy(pending["items"]), wrapper=wrapper)
run_state._current_turn_persisted_item_count = pending["persisted_count"]
# Keep the checkpoint until compaction also settles: if _apply_post_write_compaction
# raises below, a later retry must still be able to redo just the compaction step
# instead of silently losing it. The append itself is retry-safe (the reconciliation
# above detects an already-committed batch and skips re-appending it).
await _apply_post_write_compaction(
session,
response_id=pending.get("response_id"),
store=pending.get("store"),
has_local_tool_outputs=pending.get("has_local_tool_outputs", False),
wrapper=wrapper,
)
run_state._pending_session_write = None
finally:
run_state._session_write_in_progress = False
Expand Down
Loading
Loading