Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
11e8d33
fix(sessions): reset compaction state after pop
Sep 4, 2026
4044308
fix(sessions): guard response chain across pop races
Sep 5, 2026
7b0179c
fix(sessions): avoid stale compaction candidate cache
Sep 5, 2026
086b845
fix(sessions): revalidate candidate cache after mutations
Sep 5, 2026
708cdd9
fix(sessions): invalidate stale response-id compactions
Sep 5, 2026
2a9d945
test(sessions): satisfy full pyright checks
fscfede-beep Sep 5, 2026
ff5e264
fix(sessions): serialize response ownership with replacement
fscfede-beep Sep 5, 2026
dabfa13
fix(sessions): reject stale compaction ownership
fscfede-beep Sep 5, 2026
61572f6
fix(sessions): reject response id invalidated while waiting
fscfede-beep Sep 5, 2026
fd3b0ee
fix(sessions): invalidate response chain after ambiguous clear
fscfede-beep Sep 5, 2026
5702cf2
fix(sessions): invalidate chain after clear
fscfede-beep Sep 5, 2026
84508a4
merge: sync upstream main for #4868
fscfede-beep Sep 6, 2026
354d51a
fix(sessions): preserve queued response id updates
fscfede-beep Sep 7, 2026
65905d5
fix(sessions): preserve deferred compaction retry
fscfede-beep Sep 7, 2026
26fa2ad
fix(sessions): preserve deferred response ownership
fscfede-beep Sep 7, 2026
74139eb
fix(sessions): avoid resurrecting completed compaction retry
fscfede-beep Sep 7, 2026
ee829d0
fix(sessions): retain deferred retry on compact failure
fscfede-beep Sep 7, 2026
9d654fe
fix(sessions): preserve deferred retry ownership
fscfede-beep Sep 7, 2026
e172f74
Merge upstream/main into fix/compaction-pop-response-state
fscfede-beep Sep 7, 2026
141d167
fix(sessions): prevent stale deferred cache publication
fscfede-beep Sep 7, 2026
2a18ce3
fix(sessions): reject stale manual compaction
fscfede-beep Sep 7, 2026
45c3e25
fix(sessions): invalidate stale manual response chain
fscfede-beep Sep 7, 2026
7e67348
Merge branch 'main' of https://github.com/openai/openai-agents-python…
fscfede-beep Sep 7, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
129 changes: 104 additions & 25 deletions src/agents/memory/openai_responses_compaction_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,7 @@ def __init__(
self._response_id: str | None = None
self._deferred_response_id: str | None = None
self._last_unstored_response_id: str | None = None
self._response_chain_invalidation_generation = 0
# Serialize wrapper mutations against compaction snapshot/replace/restore so a
# cancellation rollback cannot rewrite past a newer concurrent write.
self._mutation_lock = asyncio.Lock()
Expand All @@ -147,6 +148,12 @@ def __init__(
# pending automatic replacement without inferring ownership from history.
self._mutation_generation = 0

def _invalidate_response_chain(self) -> None:
self._response_id = None
self._deferred_response_id = None
self._last_unstored_response_id = None
self._response_chain_invalidation_generation += 1

@property
def client(self) -> AsyncOpenAI:
if self._client is None:
Expand Down Expand Up @@ -182,10 +189,14 @@ async def run_compaction(
When a run context is provided, the billed compaction request contributes to
that run's usage totals.
"""
# Keep one wrapper mutation boundary from the snapshot through replacement.
# A concurrent add, pop, or clear waits here and then runs against the
# compacted state instead of being overwritten by a stale replacement.
# Reject a caller-supplied retained chain if a destructive mutation completed
# while this call was waiting to acquire compaction ownership.
pre_lock_invalidation_generation = self._response_chain_invalidation_generation
pre_lock_mutation_generation = self._mutation_generation
async with self._mutation_lock:
if pre_lock_invalidation_generation != self._response_chain_invalidation_generation:
logger.debug("skip: response chain invalidated while waiting for mutation lock")
return
has_expected_generation = wrapper is not None and hasattr(
wrapper, "_session_compaction_generation"
)
Expand All @@ -203,6 +214,29 @@ async def run_compaction(
"run appended its items."
)
return
if (
not has_expected_generation
and pre_lock_mutation_generation != self._mutation_generation
):
Comment thread
fscfede-beep marked this conversation as resolved.
guard_response_id = (
args.get("response_id")
if args and args.get("response_id")
else self._response_id
)
guard_store = args.get("store") if args and "store" in args else None
guard_requested_mode = args.get("compaction_mode") if args else None
guard_mode = self._resolve_compaction_mode_for_response(
response_id=guard_response_id,
store=guard_store,
requested_mode=guard_requested_mode,
)
if guard_mode == "previous_response_id":
self._invalidate_response_chain()
logger.warning(
"Skipped compaction because Session history changed while this "
"manual response-chain compaction waited for mutation ownership."
)
return
await self._run_compaction_locked(args, wrapper=wrapper)

async def _run_compaction_locked(
Expand Down Expand Up @@ -254,6 +288,7 @@ async def _run_compaction_locked(
)
return

deferred_response_id = self._deferred_response_id
self._deferred_response_id = None
logger.debug(
"compact: start for %s using %s (mode=%s)",
Expand All @@ -268,7 +303,12 @@ async def _run_compaction_locked(
else:
compact_kwargs["input"] = session_items

compacted = await self.client.responses.compact(**compact_kwargs)
try:
compacted = await self.client.responses.compact(**compact_kwargs)
except (Exception, asyncio.CancelledError):
if deferred_response_id is not None and self._deferred_response_id is None:
self._deferred_response_id = deferred_response_id
raise

compacted_usage = getattr(compacted, "usage", None)
if wrapper is not None and compacted_usage is not None:
Expand Down Expand Up @@ -429,24 +469,56 @@ async def _restore_underlying_session_items(
)

async def _defer_compaction(self, response_id: str, store: bool | None = None) -> None:
if self._deferred_response_id is not None:
return
compaction_candidate_items, session_items = await self._ensure_compaction_candidates()
resolved_mode = self._resolve_compaction_mode_for_response(
response_id=response_id,
store=store,
requested_mode=None,
)
should_compact = self.should_trigger_compaction(
{
"response_id": response_id,
"compaction_mode": resolved_mode,
"compaction_candidate_items": compaction_candidate_items,
"session_items": session_items,
}
)
if should_compact:
self._deferred_response_id = response_id
pre_lock_invalidation_generation = self._response_chain_invalidation_generation
while True:
async with self._mutation_lock:
if pre_lock_invalidation_generation != self._response_chain_invalidation_generation:
return
if self._deferred_response_id is not None:
return
invalidation_generation = self._response_chain_invalidation_generation
mutation_generation = self._mutation_generation
if self._compaction_candidate_items is not None and self._session_items is not None:
compaction_candidate_items = self._compaction_candidate_items[:]
session_items = self._session_items[:]
should_compact = self.should_trigger_compaction(
{
"response_id": response_id,
"compaction_mode": self._resolve_compaction_mode_for_response(
response_id=response_id, store=store, requested_mode=None
),
"compaction_candidate_items": compaction_candidate_items,
"session_items": session_items,
}
)
if should_compact:
self._deferred_response_id = response_id
return

history = _normalize_compaction_session_items(await self.underlying_session.get_items())
compaction_candidate_items = select_compaction_candidate_items(history)

async with self._mutation_lock:
if invalidation_generation != self._response_chain_invalidation_generation:
return
if self._deferred_response_id is not None:
return
if mutation_generation != self._mutation_generation:
pre_lock_invalidation_generation = self._response_chain_invalidation_generation
continue
should_compact = self.should_trigger_compaction(
{
"response_id": response_id,
"compaction_mode": self._resolve_compaction_mode_for_response(
response_id=response_id, store=store, requested_mode=None
),
"compaction_candidate_items": compaction_candidate_items,
"session_items": history,
}
)
if should_compact:
self._deferred_response_id = response_id
return

def _get_deferred_compaction_response_id(self) -> str | None:
return self._deferred_response_id
Expand Down Expand Up @@ -493,7 +565,13 @@ async def pop_item(self) -> TResponseInputItem | None:
async with self._mutation_lock:
try:
popped = await self.underlying_session.pop_item()
except (Exception, asyncio.CancelledError):
except asyncio.CancelledError:
self._compaction_candidate_items = None
self._session_items = None
self._mutation_generation += 1
self._invalidate_response_chain()
raise
except Exception:
self._compaction_candidate_items = None
self._session_items = None
self._mutation_generation += 1
Expand All @@ -502,6 +580,7 @@ async def pop_item(self) -> TResponseInputItem | None:
self._compaction_candidate_items = None
self._session_items = None
self._mutation_generation += 1
self._invalidate_response_chain()
return popped

async def clear_session(self) -> None:
Expand All @@ -511,13 +590,13 @@ async def clear_session(self) -> None:
except (Exception, asyncio.CancelledError):
self._compaction_candidate_items = None
self._session_items = None
self._deferred_response_id = None
self._mutation_generation += 1
self._invalidate_response_chain()
raise
self._compaction_candidate_items = []
self._session_items = []
self._deferred_response_id = None
self._mutation_generation += 1
self._invalidate_response_chain()

async def _ensure_compaction_candidates(
self,
Expand Down
Loading