From 3a642c2e252a613786502b120d5fe3e3ba1598d9 Mon Sep 17 00:00:00 2001 From: ScarabSystems Date: Tue, 28 Jul 2026 22:06:10 +0000 Subject: [PATCH 1/4] Bound summarization input before provider call SummarizationStrategy now selects complete message groups that fit a configurable summary input token budget before calling the summary client. Only messages actually sent to the summarizer are annotated and excluded, leaving oversized later groups for a later compaction pass instead of shipping the whole transcript unbounded. Validation: uv run pytest packages/core/tests/core/test_compaction.py -k bounds_summary_input -m "not integration" failed before the implementation and passed after it; uv run pytest packages/core/tests/core/test_compaction.py -m "not integration" passed; uv run poe test -P core passed; uv run poe install completed; uv run poe check -P core passed. --- .../core/agent_framework/_compaction.py | 65 +++++++++++++++++-- .../core/tests/core/test_compaction.py | 59 +++++++++++++++++ 2 files changed, 119 insertions(+), 5 deletions(-) diff --git a/python/packages/core/agent_framework/_compaction.py b/python/packages/core/agent_framework/_compaction.py index 59abb10a46..c18b9f2765 100644 --- a/python/packages/core/agent_framework/_compaction.py +++ b/python/packages/core/agent_framework/_compaction.py @@ -979,6 +979,35 @@ def _format_messages_for_summary(messages: list[Message]) -> str: return "\n".join(lines) +def _select_summary_input_groups( + groups: Sequence[tuple[str, list[Message]]], + *, + prompt: str, + max_summary_input_tokens: int | None, + tokenizer: TokenizerProtocol, +) -> tuple[list[str], list[Message]]: + if max_summary_input_tokens is None: + return ( + [group_id for group_id, _ in groups], + [message for _, group_messages in groups for message in group_messages], + ) + + selected_group_ids: list[str] = [] + selected_messages: list[Message] = [] + prompt_token_count = tokenizer.count_tokens(prompt) + + for group_id, group_messages in groups: + candidate_messages = [*selected_messages, *group_messages] + candidate_text = _format_messages_for_summary(candidate_messages) + candidate_token_count = prompt_token_count + tokenizer.count_tokens(candidate_text) + if candidate_token_count > max_summary_input_tokens: + break + selected_group_ids.append(group_id) + selected_messages = candidate_messages + + return selected_group_ids, selected_messages + + DEFAULT_SUMMARIZATION_PROMPT: Final[ str ] = """**Generate a clear and complete summary of the entire conversation in no more than five sentences.** @@ -996,6 +1025,8 @@ def _format_messages_for_summary(messages: list[Message]) -> str: - Omit any details included in an earlier summary """ +DEFAULT_SUMMARY_INPUT_TOKEN_BUDGET: Final[int] = 8_000 + class SummarizationStrategy: """Summarize older included groups and replace them with linked summary text. @@ -1026,6 +1057,8 @@ def __init__( target_count: int = 4, threshold: int | None = 2, prompt: str | None = None, + max_summary_input_tokens: int | None = DEFAULT_SUMMARY_INPUT_TOKEN_BUDGET, + tokenizer: TokenizerProtocol | None = None, ) -> None: """Create a summarization strategy. @@ -1043,19 +1076,31 @@ def __init__( prompt: Optional summarization instruction. If omitted, a default prompt that preserves goals, decisions, and unresolved items is used. + max_summary_input_tokens: Maximum estimated token count for the + summarizer request prompt and user transcript. Whole message + groups are selected until the next group would exceed this + budget. Pass ``None`` to disable the input budget. + tokenizer: Token counter used to estimate summarizer request size. + If omitted, :class:`CharacterEstimatorTokenizer` is used. Raises: ValueError: If ``target_count`` is less than 1. ValueError: If ``threshold`` is provided and is negative. + ValueError: If ``max_summary_input_tokens`` is provided and is less + than 1. """ if target_count <= 0: raise ValueError("target_count must be greater than 0.") if threshold is not None and threshold < 0: raise ValueError("threshold must be greater than or equal to 0.") + if max_summary_input_tokens is not None and max_summary_input_tokens <= 0: + raise ValueError("max_summary_input_tokens must be greater than 0.") self.client = client self.target_count = target_count self.threshold = threshold if threshold is not None else 0 self.prompt = prompt or DEFAULT_SUMMARIZATION_PROMPT + self.max_summary_input_tokens = max_summary_input_tokens + self.tokenizer = tokenizer or CharacterEstimatorTokenizer() async def __call__(self, messages: list[Message]) -> bool: ordered_group_ids = _ordered_group_ids_from_annotations(messages) @@ -1096,12 +1141,22 @@ async def __call__(self, messages: list[Message]) -> bool: if not group_ids_to_summarize: return False - messages_to_summarize: list[Message] = [] - for group_id, group_messages in included_non_system_groups: - if group_id in keep_group_id_set: - continue - messages_to_summarize.extend(group_messages) + candidate_groups = [ + (group_id, group_messages) + for group_id, group_messages in included_non_system_groups + if group_id not in keep_group_id_set + ] + group_ids_to_summarize, messages_to_summarize = _select_summary_input_groups( + candidate_groups, + prompt=self.prompt, + max_summary_input_tokens=self.max_summary_input_tokens, + tokenizer=self.tokenizer, + ) if not messages_to_summarize: + if self.max_summary_input_tokens is not None: + logger.warning( + "Skipping summarization compaction: no complete message group fits within max_summary_input_tokens." + ) return False try: diff --git a/python/packages/core/tests/core/test_compaction.py b/python/packages/core/tests/core/test_compaction.py index 4507d39946..2225a0bad5 100644 --- a/python/packages/core/tests/core/test_compaction.py +++ b/python/packages/core/tests/core/test_compaction.py @@ -477,6 +477,27 @@ async def get_response( return ChatResponse(messages=[Message(role="assistant", contents=[" "])]) +class _RecordingSummarizer: + def __init__(self) -> None: + self.requests: list[list[Message]] = [] + + async def get_response( + self, + messages: list[Message], + *, + stream: bool = False, + options: dict[str, Any] | None = None, + **kwargs: Any, + ) -> ChatResponse: + self.requests.append(messages) + return ChatResponse(messages=[Message(role="assistant", contents=["budgeted summary"])]) + + +class _CharacterCountTokenizer: + def count_tokens(self, text: str) -> int: + return len(text) + + async def test_summarization_strategy_adds_bidirectional_trace_links() -> None: messages = [ Message(role="user", contents=["u1"]), @@ -508,6 +529,44 @@ async def test_summarization_strategy_adds_bidirectional_trace_links() -> None: assert message.additional_properties.get(EXCLUDED_KEY) is True +async def test_summarization_strategy_bounds_summary_input_to_complete_groups() -> None: + summarizer = _RecordingSummarizer() + messages = [ + Message(role="user", contents=["first old " * 20]), + Message(role="assistant", contents=["second oversized " * 120]), + Message(role="user", contents=["third should wait"]), + Message(role="assistant", contents=["fourth should wait"]), + Message(role="user", contents=["recent user"]), + Message(role="assistant", contents=["recent assistant"]), + ] + first_old_message = messages[0] + oversized_message = messages[1] + strategy = SummarizationStrategy( + client=summarizer, # type: ignore[arg-type] # pyrefly: ignore[bad-argument-type] # ty: ignore[invalid-argument-type] + target_count=2, + threshold=0, + max_summary_input_tokens=1_000, + tokenizer=_CharacterCountTokenizer(), + ) + annotate_message_groups(messages) + + changed = await strategy(messages) + + assert changed is True + assert len(summarizer.requests) == 1 + summary_request_text = summarizer.requests[0][1].text + assert summary_request_text is not None + assert "first old" in summary_request_text + assert "second oversized" not in summary_request_text + assert first_old_message.additional_properties.get(EXCLUDED_KEY) is True + assert oversized_message.additional_properties.get(EXCLUDED_KEY) is not True + summary = next(message for message in messages if _group_unknown_value(message, SUMMARY_OF_MESSAGE_IDS_KEY)) + summarized_message_ids = _group_unknown_value(summary, SUMMARY_OF_MESSAGE_IDS_KEY) + assert isinstance(summarized_message_ids, list) + assert first_old_message.message_id in summarized_message_ids + assert oversized_message.message_id not in summarized_message_ids + + async def test_summarization_strategy_returns_false_when_summary_generation_fails( caplog: Any, ) -> None: From 479d97520b3e9ac91f6f464a5a29ebe66c058c82 Mon Sep 17 00:00:00 2001 From: ScarabSystems Date: Tue, 28 Jul 2026 18:26:58 -0400 Subject: [PATCH 2/4] Handle oversized leading summary groups Skip individually over-budget leading groups when selecting summarization input so a large early transcript item does not prevent later compactable groups from being summarized. Validation: uv run pytest packages/core/tests/core/test_compaction.py -k skips_oversized_first_group -q; uv run pytest packages/core/tests/core/test_compaction.py -q; uv run poe check -P core. --- .../core/agent_framework/_compaction.py | 2 ++ .../core/tests/core/test_compaction.py | 31 +++++++++++++++++++ 2 files changed, 33 insertions(+) diff --git a/python/packages/core/agent_framework/_compaction.py b/python/packages/core/agent_framework/_compaction.py index c18b9f2765..d574b9d9ff 100644 --- a/python/packages/core/agent_framework/_compaction.py +++ b/python/packages/core/agent_framework/_compaction.py @@ -1001,6 +1001,8 @@ def _select_summary_input_groups( candidate_text = _format_messages_for_summary(candidate_messages) candidate_token_count = prompt_token_count + tokenizer.count_tokens(candidate_text) if candidate_token_count > max_summary_input_tokens: + if not selected_messages: + continue break selected_group_ids.append(group_id) selected_messages = candidate_messages diff --git a/python/packages/core/tests/core/test_compaction.py b/python/packages/core/tests/core/test_compaction.py index 2225a0bad5..c17aa794d1 100644 --- a/python/packages/core/tests/core/test_compaction.py +++ b/python/packages/core/tests/core/test_compaction.py @@ -567,6 +567,37 @@ async def test_summarization_strategy_bounds_summary_input_to_complete_groups() assert oversized_message.message_id not in summarized_message_ids +async def test_summarization_strategy_skips_oversized_first_group() -> None: + summarizer = _RecordingSummarizer() + messages = [ + Message(role="user", contents=["oversized first group " * 120]), + Message(role="assistant", contents=["small later group"]), + Message(role="user", contents=["recent user"]), + Message(role="assistant", contents=["recent assistant"]), + ] + oversized_message = messages[0] + small_message = messages[1] + strategy = SummarizationStrategy( + client=summarizer, # type: ignore[arg-type] # pyrefly: ignore[bad-argument-type] # ty: ignore[invalid-argument-type] + target_count=2, + threshold=0, + max_summary_input_tokens=1_000, + tokenizer=_CharacterCountTokenizer(), + ) + annotate_message_groups(messages) + + changed = await strategy(messages) + + assert changed is True + assert len(summarizer.requests) == 1 + summary_request_text = summarizer.requests[0][1].text + assert summary_request_text is not None + assert "oversized first group" not in summary_request_text + assert "small later group" in summary_request_text + assert oversized_message.additional_properties.get(EXCLUDED_KEY) is not True + assert small_message.additional_properties.get(EXCLUDED_KEY) is True + + async def test_summarization_strategy_returns_false_when_summary_generation_fails( caplog: Any, ) -> None: From ae6269ccc0d1bea67e2e5c7879bf745084c8f4ae Mon Sep 17 00:00:00 2001 From: ScarabSystems Date: Tue, 28 Jul 2026 18:50:12 -0400 Subject: [PATCH 3/4] Escalate repeated summary failures Track consecutive SummarizationStrategy failures and emit a single error once the strategy has failed three times without a successful summary. Reset the escalation state after a successful summary so only persistent failures become loud. Validation: uv run pytest packages/core/tests/core/test_compaction.py -k 'repeated_summary_failures or resets_failure_escalation' -q; uv run pytest packages/core/tests/core/test_compaction.py -q; uv run poe check -P core. --- .../core/agent_framework/_compaction.py | 24 +++++++ .../core/tests/core/test_compaction.py | 69 +++++++++++++++++++ 2 files changed, 93 insertions(+) diff --git a/python/packages/core/agent_framework/_compaction.py b/python/packages/core/agent_framework/_compaction.py index d574b9d9ff..4a76b81ce8 100644 --- a/python/packages/core/agent_framework/_compaction.py +++ b/python/packages/core/agent_framework/_compaction.py @@ -1028,6 +1028,7 @@ def _select_summary_input_groups( """ DEFAULT_SUMMARY_INPUT_TOKEN_BUDGET: Final[int] = 8_000 +SUMMARY_FAILURE_ERROR_THRESHOLD: Final[int] = 3 class SummarizationStrategy: @@ -1103,6 +1104,25 @@ def __init__( self.prompt = prompt or DEFAULT_SUMMARIZATION_PROMPT self.max_summary_input_tokens = max_summary_input_tokens self.tokenizer = tokenizer or CharacterEstimatorTokenizer() + self._consecutive_summary_failures = 0 + self._summary_failure_error_emitted = False + + def _record_summary_failure(self) -> None: + self._consecutive_summary_failures += 1 + if ( + self._consecutive_summary_failures >= SUMMARY_FAILURE_ERROR_THRESHOLD + and not self._summary_failure_error_emitted + ): + logger.error( + "Summarization compaction has failed %s consecutive times; " + "graceful summary compaction may no longer be contributing.", + self._consecutive_summary_failures, + ) + self._summary_failure_error_emitted = True + + def _record_summary_success(self) -> None: + self._consecutive_summary_failures = 0 + self._summary_failure_error_emitted = False async def __call__(self, messages: list[Message]) -> bool: ordered_group_ids = _ordered_group_ids_from_annotations(messages) @@ -1159,6 +1179,7 @@ async def __call__(self, messages: list[Message]) -> bool: logger.warning( "Skipping summarization compaction: no complete message group fits within max_summary_input_tokens." ) + self._record_summary_failure() return False try: @@ -1177,12 +1198,15 @@ async def __call__(self, messages: list[Message]) -> bool: "Skipping summarization compaction: summary generation failed (%s).", exc, ) + self._record_summary_failure() return False summary_text = summary_response.text.strip() if summary_response.text else "" if not summary_text: logger.warning("Skipping summarization compaction: summarizer returned no text.") + self._record_summary_failure() return False + self._record_summary_success() summary_id = f"summary_{len(messages)}" original_message_ids = [message.message_id for message in messages_to_summarize if message.message_id] summary_of_group_ids = list(group_ids_to_summarize) diff --git a/python/packages/core/tests/core/test_compaction.py b/python/packages/core/tests/core/test_compaction.py index c17aa794d1..efd0543a8b 100644 --- a/python/packages/core/tests/core/test_compaction.py +++ b/python/packages/core/tests/core/test_compaction.py @@ -477,6 +477,24 @@ async def get_response( return ChatResponse(messages=[Message(role="assistant", contents=[" "])]) +class _ScriptedSummarizer: + def __init__(self, outcomes: list[str | BaseException]) -> None: + self.outcomes = outcomes + + async def get_response( + self, + messages: list[Message], + *, + stream: bool = False, + options: dict[str, Any] | None = None, + **kwargs: Any, + ) -> ChatResponse: + outcome = self.outcomes.pop(0) + if isinstance(outcome, BaseException): + raise outcome + return ChatResponse(messages=[Message(role="assistant", contents=[outcome])]) + + class _RecordingSummarizer: def __init__(self) -> None: self.requests: list[list[Message]] = [] @@ -620,6 +638,57 @@ async def test_summarization_strategy_returns_false_when_summary_generation_fail assert all(message.additional_properties.get(EXCLUDED_KEY) is not True for message in messages) +async def test_summarization_strategy_escalates_repeated_summary_failures(caplog: Any) -> None: + messages = [ + Message(role="user", contents=["u1"]), + Message(role="assistant", contents=["a1"]), + Message(role="user", contents=["u2"]), + Message(role="assistant", contents=["a2"]), + Message(role="user", contents=["u3"]), + Message(role="assistant", contents=["a3"]), + ] + strategy = SummarizationStrategy(client=_FailingSummarizer(), target_count=2, threshold=0) # type: ignore[arg-type] # pyrefly: ignore[bad-argument-type] # ty: ignore[invalid-argument-type] + annotate_message_groups(messages) + + with caplog.at_level(logging.WARNING, logger="agent_framework"): + assert await strategy(messages) is False + assert await strategy(messages) is False + assert await strategy(messages) is False + assert await strategy(messages) is False + + error_records = [record for record in caplog.records if record.levelno == logging.ERROR] + assert len(error_records) == 1 + assert "failed 3 consecutive times" in error_records[0].message + + +async def test_summarization_strategy_resets_failure_escalation_after_success( + caplog: Any, +) -> None: + summarizer = _ScriptedSummarizer([ + RuntimeError("first failure"), + RuntimeError("second failure"), + "recovered summary", + RuntimeError("third failure"), + RuntimeError("fourth failure"), + ]) + strategy = SummarizationStrategy(client=summarizer, target_count=2, threshold=0) # type: ignore[arg-type] # pyrefly: ignore[bad-argument-type] # ty: ignore[invalid-argument-type] + + with caplog.at_level(logging.WARNING, logger="agent_framework"): + for _ in range(5): + messages = [ + Message(role="user", contents=["u1"]), + Message(role="assistant", contents=["a1"]), + Message(role="user", contents=["u2"]), + Message(role="assistant", contents=["a2"]), + Message(role="user", contents=["u3"]), + Message(role="assistant", contents=["a3"]), + ] + annotate_message_groups(messages) + await strategy(messages) + + assert not any(record.levelno == logging.ERROR for record in caplog.records) + + async def test_summarization_strategy_returns_false_when_summary_is_empty( caplog: Any, ) -> None: From 086442c0df301bcdc95dd6f621280938ac5c0eb0 Mon Sep 17 00:00:00 2001 From: ScarabSystems Date: Tue, 28 Jul 2026 20:48:38 -0400 Subject: [PATCH 4/4] Refine summary input selection Avoid rebuilding and re-tokenizing the full selected summary transcript on every candidate group while preserving complete-group selection and oversized leading group skipping. Tighten the scripted summarizer test helper to expected Exception failures instead of BaseException. Verification: uv run pytest packages/core/tests/core/test_compaction.py -q; uv run poe syntax -P core. --- .../core/agent_framework/_compaction.py | 31 +++++++++----- .../core/tests/core/test_compaction.py | 41 ++++++++++++++++++- 2 files changed, 60 insertions(+), 12 deletions(-) diff --git a/python/packages/core/agent_framework/_compaction.py b/python/packages/core/agent_framework/_compaction.py index 4a76b81ce8..9b81d96158 100644 --- a/python/packages/core/agent_framework/_compaction.py +++ b/python/packages/core/agent_framework/_compaction.py @@ -969,13 +969,17 @@ def _tool_result_text(value: Any) -> str: return str(cast(object, value)) -def _format_messages_for_summary(messages: list[Message]) -> str: +def _format_summary_message(index: int, message: Message) -> str: + content_text = message.text + if not content_text: + content_text = ", ".join(content.type for content in message.contents) + return f"{index}. [{message.role}] {content_text}" + + +def _format_messages_for_summary(messages: list[Message], *, start_index: int = 1) -> str: lines: list[str] = [] - for index, message in enumerate(messages, start=1): - content_text = message.text - if not content_text: - content_text = ", ".join(content.type for content in message.contents) - lines.append(f"{index}. [{message.role}] {content_text}") + for index, message in enumerate(messages, start=start_index): + lines.append(_format_summary_message(index, message)) return "\n".join(lines) @@ -995,17 +999,24 @@ def _select_summary_input_groups( selected_group_ids: list[str] = [] selected_messages: list[Message] = [] prompt_token_count = tokenizer.count_tokens(prompt) + selected_message_count = 0 + selected_text_token_count = 0 + separator_token_count = tokenizer.count_tokens("\n") for group_id, group_messages in groups: - candidate_messages = [*selected_messages, *group_messages] - candidate_text = _format_messages_for_summary(candidate_messages) - candidate_token_count = prompt_token_count + tokenizer.count_tokens(candidate_text) + group_text = _format_messages_for_summary(group_messages, start_index=selected_message_count + 1) + candidate_text_token_count = selected_text_token_count + tokenizer.count_tokens(group_text) + if selected_messages: + candidate_text_token_count += separator_token_count + candidate_token_count = prompt_token_count + candidate_text_token_count if candidate_token_count > max_summary_input_tokens: if not selected_messages: continue break selected_group_ids.append(group_id) - selected_messages = candidate_messages + selected_messages.extend(group_messages) + selected_message_count += len(group_messages) + selected_text_token_count = candidate_text_token_count return selected_group_ids, selected_messages diff --git a/python/packages/core/tests/core/test_compaction.py b/python/packages/core/tests/core/test_compaction.py index efd0543a8b..e4c1870731 100644 --- a/python/packages/core/tests/core/test_compaction.py +++ b/python/packages/core/tests/core/test_compaction.py @@ -35,6 +35,7 @@ included_token_count, ) from agent_framework._compaction import ( + _select_summary_input_groups, _serialize_message, append_compaction_message, extend_compaction_messages, @@ -478,7 +479,7 @@ async def get_response( class _ScriptedSummarizer: - def __init__(self, outcomes: list[str | BaseException]) -> None: + def __init__(self, outcomes: list[str | Exception]) -> None: self.outcomes = outcomes async def get_response( @@ -490,7 +491,7 @@ async def get_response( **kwargs: Any, ) -> ChatResponse: outcome = self.outcomes.pop(0) - if isinstance(outcome, BaseException): + if isinstance(outcome, Exception): raise outcome return ChatResponse(messages=[Message(role="assistant", contents=[outcome])]) @@ -516,6 +517,15 @@ def count_tokens(self, text: str) -> int: return len(text) +class _RecordingCharacterCountTokenizer: + def __init__(self) -> None: + self.seen_texts: list[str] = [] + + def count_tokens(self, text: str) -> int: + self.seen_texts.append(text) + return len(text) + + async def test_summarization_strategy_adds_bidirectional_trace_links() -> None: messages = [ Message(role="user", contents=["u1"]), @@ -616,6 +626,33 @@ async def test_summarization_strategy_skips_oversized_first_group() -> None: assert small_message.additional_properties.get(EXCLUDED_KEY) is True +def test_summary_input_selection_does_not_retokenize_selected_transcript() -> None: + tokenizer = _RecordingCharacterCountTokenizer() + groups = [ + ("group_1", [Message(role="user", contents=["first"])]), + ("group_2", [Message(role="assistant", contents=["second"])]), + ("group_3", [Message(role="user", contents=["third"])]), + ] + + selected_group_ids, selected_messages = _select_summary_input_groups( + groups, + prompt="prompt", + max_summary_input_tokens=1_000, + tokenizer=tokenizer, + ) + + assert selected_group_ids == ["group_1", "group_2", "group_3"] + assert selected_messages == [message for _, group_messages in groups for message in group_messages] + assert ( + "\n".join([ + "1. [user] first", + "2. [assistant] second", + "3. [user] third", + ]) + not in tokenizer.seen_texts + ) + + async def test_summarization_strategy_returns_false_when_summary_generation_fails( caplog: Any, ) -> None: