From 929ec269f05610850ecf79564bc481df5eded8e0 Mon Sep 17 00:00:00 2001 From: nanhanq1 <157583847+nanhanq1@users.noreply.github.com> Date: Sun, 27 Sep 2026 20:13:45 +0800 Subject: [PATCH 1/2] fix: release the token counter when CompactionHook is closed --- haystack/hooks/compaction/hooks.py | 51 +++++++++++++------ ...-close-token-counter-0495043806039056.yaml | 10 ++++ test/hooks/compaction/test_hooks.py | 49 +++++++++++++++++- 3 files changed, 93 insertions(+), 17 deletions(-) create mode 100644 releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml diff --git a/haystack/hooks/compaction/hooks.py b/haystack/hooks/compaction/hooks.py index cc1cc124a4d..8b977cfb917 100644 --- a/haystack/hooks/compaction/hooks.py +++ b/haystack/hooks/compaction/hooks.py @@ -50,6 +50,32 @@ def _estimated_context_tokens( return context_tokens + token_counter.count(messages=tool_result_messages) +async def _warm_up_async(resource: Any) -> None: + """ + Warm up one lifecycle-bearing resource, awaiting `warm_up_async` when defined. + + :param resource: The token counter or compactor to warm up. + """ + warm_up_async = getattr(resource, "warm_up_async", None) + if warm_up_async is not None: + await warm_up_async() + elif hasattr(resource, "warm_up"): + resource.warm_up() + + +async def _close_async(resource: Any) -> None: + """ + Release one lifecycle-bearing resource, awaiting `close_async` when defined. + + :param resource: The token counter or compactor to close. + """ + close_async = getattr(resource, "close_async", None) + if close_async is not None: + await close_async() + elif hasattr(resource, "close"): + resource.close() + + @_experimental class CompactionHook: """ @@ -288,26 +314,19 @@ def warm_up(self) -> None: async def warm_up_async(self) -> None: """Warm up the token counter and the compactor on the serving event loop.""" - if hasattr(self.token_counter, "warm_up"): - self.token_counter.warm_up() - warm_up_async = getattr(self.compactor, "warm_up_async", None) - if warm_up_async is not None: - await warm_up_async() - elif hasattr(self.compactor, "warm_up"): - self.compactor.warm_up() + await _warm_up_async(resource=self.token_counter) + await _warm_up_async(resource=self.compactor) def close(self) -> None: - """Release the compactor's resources.""" - if hasattr(self.compactor, "close"): - self.compactor.close() + """Release the token counter's and the compactor's resources.""" + for resource in (self.token_counter, self.compactor): + if hasattr(resource, "close"): + resource.close() async def close_async(self) -> None: - """Release the compactor's async resources.""" - close_async = getattr(self.compactor, "close_async", None) - if close_async is not None: - await close_async() - elif hasattr(self.compactor, "close"): - self.compactor.close() + """Release the token counter's and the compactor's async resources.""" + for resource in (self.token_counter, self.compactor): + await _close_async(resource=resource) def to_dict(self) -> dict[str, Any]: """ diff --git a/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml b/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml new file mode 100644 index 00000000000..9355d4553c0 --- /dev/null +++ b/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml @@ -0,0 +1,10 @@ +--- +fixes: + - | + Fixed ``CompactionHook`` never releasing its ``token_counter``. ``warm_up()`` and + ``warm_up_async()`` have always warmed up both the token counter and the compactor, but + ``close()`` and ``close_async()`` released only the compactor, so a counter that holds a + resource stayed open after the ``Agent`` was closed. ``OpenAITokenCounter`` is the counter + where this is visible: it builds an ``OpenAI`` client, and with it an HTTP connection pool, + in ``warm_up()`` and drops it again in its own ``close()``, which nothing ever called. Both + lifecycle methods now handle the token counter exactly as they already handle the compactor. diff --git a/test/hooks/compaction/test_hooks.py b/test/hooks/compaction/test_hooks.py index e32f6458068..9b731c6e665 100644 --- a/test/hooks/compaction/test_hooks.py +++ b/test/hooks/compaction/test_hooks.py @@ -16,7 +16,7 @@ from haystack.hooks.compaction.hooks import _estimated_context_tokens from haystack.hooks.compaction.utils import _COMPACTION_META_KEY, _last_assistant_index from haystack.hooks.invocation import _run_hooks, _run_hooks_async -from haystack.token_counters import TokenCounter +from haystack.token_counters import OpenAITokenCounter, TokenCounter from haystack.tools import tool from haystack.utils.experimental import ExperimentalWarning from test.hooks.compaction.helpers import ( @@ -87,6 +87,26 @@ def to_dict(self) -> dict[str, Any]: return default_to_dict(self) +class _RecordingCounter(FakeCounter): + """A counter that records the lifecycle calls made to it, like `OpenAITokenCounter` does.""" + + def __init__(self) -> None: + super().__init__() + self.calls: list[str] = [] + + def warm_up(self) -> None: + self.calls.append("warm_up") + + async def warm_up_async(self) -> None: + self.calls.append("warm_up_async") + + def close(self) -> None: + self.calls.append("close") + + async def close_async(self) -> None: + self.calls.append("close_async") + + def _hook(compactor: Compactor | None = None, **overrides: Any) -> CompactionHook: settings: dict[str, Any] = { "context_window": WINDOW, @@ -394,6 +414,25 @@ def test_lifecycle_delegates_to_the_compactor(self): hook.close() assert compactor.calls == ["warm_up", "close"] + def test_lifecycle_also_closes_the_token_counter(self): + # `warm_up()` has always warmed up both resources; `close()` used to release only the compactor. + counter = _RecordingCounter() + hook = _hook(token_counter=counter) + hook.warm_up() + hook.close() + assert counter.calls == ["warm_up", "close"] + + def test_lifecycle_closes_a_real_openai_token_counter(self, monkeypatch): + # The counter that actually holds something to release: `warm_up()` builds an `OpenAI` client and `close()` + # drops it again. Nothing in the hook used to reach that `close()`, so the client outlived the hook. + monkeypatch.setenv("OPENAI_API_KEY", "sk-test") + counter = OpenAITokenCounter(model="gpt-5-mini") + hook = _hook(token_counter=counter) + hook.warm_up() + assert counter.client is not None + hook.close() + assert counter.client is None + class TestCompactionHookInAgent: def test_compacts_a_multi_step_run(self): @@ -429,6 +468,14 @@ async def test_lifecycle_prefers_the_async_methods(self): await hook.close_async() assert compactor.calls == ["warm_up_async", "close_async"] + @pytest.mark.asyncio + async def test_lifecycle_also_closes_the_token_counter_async(self): + counter = _RecordingCounter() + hook = _hook(token_counter=counter) + await hook.warm_up_async() + await hook.close_async() + assert counter.calls == ["warm_up_async", "close_async"] + @pytest.mark.asyncio async def test_compacts_a_multi_step_async_run(self): result = await _agent({"before_llm": [_hook()]}).run_async(messages=[ChatMessage.from_user("start")]) From c67976142f45b858ee91798d775975d2e3c14f82 Mon Sep 17 00:00:00 2001 From: anakin87 Date: Mon, 28 Sep 2026 11:37:16 +0200 Subject: [PATCH 2/2] simplify --- haystack/hooks/compaction/hooks.py | 39 +++-------- ...-close-token-counter-0495043806039056.yaml | 10 +-- test/hooks/compaction/test_hooks.py | 64 ++++--------------- 3 files changed, 25 insertions(+), 88 deletions(-) diff --git a/haystack/hooks/compaction/hooks.py b/haystack/hooks/compaction/hooks.py index 8b977cfb917..d6aafecd41a 100644 --- a/haystack/hooks/compaction/hooks.py +++ b/haystack/hooks/compaction/hooks.py @@ -50,32 +50,6 @@ def _estimated_context_tokens( return context_tokens + token_counter.count(messages=tool_result_messages) -async def _warm_up_async(resource: Any) -> None: - """ - Warm up one lifecycle-bearing resource, awaiting `warm_up_async` when defined. - - :param resource: The token counter or compactor to warm up. - """ - warm_up_async = getattr(resource, "warm_up_async", None) - if warm_up_async is not None: - await warm_up_async() - elif hasattr(resource, "warm_up"): - resource.warm_up() - - -async def _close_async(resource: Any) -> None: - """ - Release one lifecycle-bearing resource, awaiting `close_async` when defined. - - :param resource: The token counter or compactor to close. - """ - close_async = getattr(resource, "close_async", None) - if close_async is not None: - await close_async() - elif hasattr(resource, "close"): - resource.close() - - @_experimental class CompactionHook: """ @@ -314,8 +288,12 @@ def warm_up(self) -> None: async def warm_up_async(self) -> None: """Warm up the token counter and the compactor on the serving event loop.""" - await _warm_up_async(resource=self.token_counter) - await _warm_up_async(resource=self.compactor) + if hasattr(self.token_counter, "warm_up"): + self.token_counter.warm_up() + if hasattr(self.compactor, "warm_up_async"): + await self.compactor.warm_up_async() + elif hasattr(self.compactor, "warm_up"): + self.compactor.warm_up() def close(self) -> None: """Release the token counter's and the compactor's resources.""" @@ -326,7 +304,10 @@ def close(self) -> None: async def close_async(self) -> None: """Release the token counter's and the compactor's async resources.""" for resource in (self.token_counter, self.compactor): - await _close_async(resource=resource) + if hasattr(resource, "close_async"): + await resource.close_async() + elif hasattr(resource, "close"): + resource.close() def to_dict(self) -> dict[str, Any]: """ diff --git a/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml b/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml index 9355d4553c0..81ba3bc8cc6 100644 --- a/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml +++ b/releasenotes/notes/compaction-hook-close-token-counter-0495043806039056.yaml @@ -1,10 +1,6 @@ --- fixes: - | - Fixed ``CompactionHook`` never releasing its ``token_counter``. ``warm_up()`` and - ``warm_up_async()`` have always warmed up both the token counter and the compactor, but - ``close()`` and ``close_async()`` released only the compactor, so a counter that holds a - resource stayed open after the ``Agent`` was closed. ``OpenAITokenCounter`` is the counter - where this is visible: it builds an ``OpenAI`` client, and with it an HTTP connection pool, - in ``warm_up()`` and drops it again in its own ``close()``, which nothing ever called. Both - lifecycle methods now handle the token counter exactly as they already handle the compactor. + Fixed ``CompactionHook.close()`` and ``close_async()`` to release the token counter's + resources as well as the compactor's. Previously, resources such as + ``OpenAITokenCounter``'s HTTP client remained open after the hook or Agent was closed. diff --git a/test/hooks/compaction/test_hooks.py b/test/hooks/compaction/test_hooks.py index 9b731c6e665..1bb8f8d5ee2 100644 --- a/test/hooks/compaction/test_hooks.py +++ b/test/hooks/compaction/test_hooks.py @@ -4,6 +4,7 @@ import logging from typing import Annotated, Any +from unittest.mock import Mock import pytest @@ -16,7 +17,7 @@ from haystack.hooks.compaction.hooks import _estimated_context_tokens from haystack.hooks.compaction.utils import _COMPACTION_META_KEY, _last_assistant_index from haystack.hooks.invocation import _run_hooks, _run_hooks_async -from haystack.token_counters import OpenAITokenCounter, TokenCounter +from haystack.token_counters import TokenCounter from haystack.tools import tool from haystack.utils.experimental import ExperimentalWarning from test.hooks.compaction.helpers import ( @@ -87,26 +88,6 @@ def to_dict(self) -> dict[str, Any]: return default_to_dict(self) -class _RecordingCounter(FakeCounter): - """A counter that records the lifecycle calls made to it, like `OpenAITokenCounter` does.""" - - def __init__(self) -> None: - super().__init__() - self.calls: list[str] = [] - - def warm_up(self) -> None: - self.calls.append("warm_up") - - async def warm_up_async(self) -> None: - self.calls.append("warm_up_async") - - def close(self) -> None: - self.calls.append("close") - - async def close_async(self) -> None: - self.calls.append("close_async") - - def _hook(compactor: Compactor | None = None, **overrides: Any) -> CompactionHook: settings: dict[str, Any] = { "context_window": WINDOW, @@ -407,32 +388,16 @@ def test_chains_tool_result_pruning_before_sliding_window(self): ) assert compacted[-2:] == messages[-2:] - def test_lifecycle_delegates_to_the_compactor(self): + def test_lifecycle_delegates_to_the_counter_and_compactor(self): + counter = Mock(spec=["warm_up", "close"]) compactor = _RecordingCompactor() - hook = _hook(compactor) + hook = _hook(compactor=compactor, token_counter=counter) hook.warm_up() hook.close() + counter.warm_up.assert_called_once_with() + counter.close.assert_called_once_with() assert compactor.calls == ["warm_up", "close"] - def test_lifecycle_also_closes_the_token_counter(self): - # `warm_up()` has always warmed up both resources; `close()` used to release only the compactor. - counter = _RecordingCounter() - hook = _hook(token_counter=counter) - hook.warm_up() - hook.close() - assert counter.calls == ["warm_up", "close"] - - def test_lifecycle_closes_a_real_openai_token_counter(self, monkeypatch): - # The counter that actually holds something to release: `warm_up()` builds an `OpenAI` client and `close()` - # drops it again. Nothing in the hook used to reach that `close()`, so the client outlived the hook. - monkeypatch.setenv("OPENAI_API_KEY", "sk-test") - counter = OpenAITokenCounter(model="gpt-5-mini") - hook = _hook(token_counter=counter) - hook.warm_up() - assert counter.client is not None - hook.close() - assert counter.client is None - class TestCompactionHookInAgent: def test_compacts_a_multi_step_run(self): @@ -461,21 +426,16 @@ async def test_run_async_uses_the_async_compaction_path(self): assert compactor.calls == ["compact_async"] @pytest.mark.asyncio - async def test_lifecycle_prefers_the_async_methods(self): + async def test_lifecycle_prefers_async_methods_with_sync_fallback(self): + counter = Mock(spec=["warm_up", "close"]) compactor = _RecordingCompactor() - hook = _hook(compactor) + hook = _hook(compactor=compactor, token_counter=counter) await hook.warm_up_async() await hook.close_async() + counter.warm_up.assert_called_once_with() + counter.close.assert_called_once_with() assert compactor.calls == ["warm_up_async", "close_async"] - @pytest.mark.asyncio - async def test_lifecycle_also_closes_the_token_counter_async(self): - counter = _RecordingCounter() - hook = _hook(token_counter=counter) - await hook.warm_up_async() - await hook.close_async() - assert counter.calls == ["warm_up_async", "close_async"] - @pytest.mark.asyncio async def test_compacts_a_multi_step_async_run(self): result = await _agent({"before_llm": [_hook()]}).run_async(messages=[ChatMessage.from_user("start")])