From 49923cbcfad28c2f765db1123775ebd5c5c33941 Mon Sep 17 00:00:00 2001 From: Mariam Amin <128838373+Mariam-Amin12@users.noreply.github.com> Date: Sun, 27 Sep 2026 16:26:37 +0300 Subject: [PATCH] feat: add ttl and time_field support to CacheChecker --- .../caching/cachechecker.mdx | 51 ++++++++- haystack/components/caching/cache_checker.py | 64 ++++++++++- .../cache-checker-ttl-6842b76bcab8efd4.yaml | 5 + test/components/caching/test_cache_checker.py | 103 +++++++++++++++++- .../caching/test_cache_checker_async.py | 23 ++++ 5 files changed, 238 insertions(+), 8 deletions(-) create mode 100644 releasenotes/notes/cache-checker-ttl-6842b76bcab8efd4.yaml diff --git a/docs-website/docs/pipeline-components/caching/cachechecker.mdx b/docs-website/docs/pipeline-components/caching/cachechecker.mdx index ae74b8efded..a8ce8e9a4b1 100644 --- a/docs-website/docs/pipeline-components/caching/cachechecker.mdx +++ b/docs-website/docs/pipeline-components/caching/cachechecker.mdx @@ -15,8 +15,9 @@ This component checks for the presence of documents in a Document Store based on | --- | --- | | **Most common position in a pipeline** | Flexible | | **Mandatory init variables** | `document_store`: A Document Store instance

`cache_field`: Name of the document's metadata field | +| **Optional init variables** | `ttl`: Time-to-live for cache entries. Can be a `timedelta` or a number of seconds. Defaults to `None`, meaning cache entries never expire

`time_field`: Name of the document's metadata field holding the cache timestamp, used only when `ttl` is set. Defaults to `"cached_at"` | | **Mandatory run variables** | `items`: A list of values associated with the `cache_field` in documents | -| **Output variables** | `hits`: A list of documents that were found with the specified value in cache

`misses`: A list of values that could not be found | +| **Output variables** | `hits`: A list of documents that were found with the specified value in cache and, if a `ttl` is configured, are still within it

`misses`: A list of values that could not be found, or whose only matches had expired | | **API reference** | [Caching](/reference/caching-api) | | **GitHub link** | https://github.com/deepset-ai/haystack/blob/main/haystack/components/caching/cache_checker.py | | **Package name** | `haystack-ai` | @@ -63,6 +64,54 @@ print( ) # Values that were not found in the cache, like: ["ABCDE"] ``` +### Expiring cache entries with a TTL + +`CacheChecker` treats a cache entry as valid forever unless you tell it otherwise. For data that goes stale — like a periodically re-crawled URL or a document whose source changes over time — pass a `ttl` so old entries stop counting as hits. + +`CacheChecker` doesn't stamp documents itself; it only reads a timestamp that was already written. Set that timestamp yourself wherever you write documents into the store: + +```python +from datetime import datetime, timezone +from haystack import Document + +document = Document( + content="doc1", + meta={ + "url": "https://example.com/resource", + "cached_at": datetime.now(timezone.utc), + }, +) +my_doc_store.write_documents([document]) +``` + +Then configure `CacheChecker` with a `ttl` so it knows how far back to look: + +```python +from datetime import timedelta +from haystack.components.caching import CacheChecker + +cache_checker = CacheChecker( + document_store=my_doc_store, + cache_field="url", + ttl=timedelta(hours=24), +) +``` + +With this configuration, `checker.run(items=["https://example.com/resource"])` only reports a hit if `cached_at` is less than 24 hours old. Anything older — or a document that never got a `cached_at` value in the first place — comes back in `misses`, exactly as if it had never been cached at all. + +The timestamp field name is configurable via `time_field` (it defaults to `"cached_at"`), which is useful if your metadata already tracks time under a different key, such as `created_at` or `last_indexed`: + +```python +cache_checker = CacheChecker( + document_store=my_doc_store, + cache_field="url", + ttl=timedelta(hours=24), + time_field="last_indexed", +) +``` + +Leaving `ttl` unset keeps the component's original behavior: cache entries never expire, and `time_field` is ignored entirely. + ### In a pipeline ```python diff --git a/haystack/components/caching/cache_checker.py b/haystack/components/caching/cache_checker.py index 2da90c84786..c2978552196 100644 --- a/haystack/components/caching/cache_checker.py +++ b/haystack/components/caching/cache_checker.py @@ -2,6 +2,7 @@ # # SPDX-License-Identifier: Apache-2.0 +from datetime import datetime, timedelta, timezone from typing import Any from haystack import Document, component, default_from_dict, default_to_dict @@ -37,7 +38,14 @@ class CacheChecker: ``` """ - def __init__(self, document_store: DocumentStore, cache_field: str) -> None: + def __init__( + self, + document_store: DocumentStore, + cache_field: str, + *, + ttl: float | timedelta | None = None, + time_field: str = "cached_at", + ) -> None: """ Creates a CacheChecker component. @@ -49,6 +57,8 @@ def __init__(self, document_store: DocumentStore, cache_field: str) -> None: """ self.document_store = document_store self.cache_field = cache_field + self.ttl = timedelta(seconds=ttl) if isinstance(ttl, (int, float)) else ttl + self.time_field = time_field def to_dict(self) -> dict[str, Any]: """ @@ -57,7 +67,46 @@ def to_dict(self) -> dict[str, Any]: :returns: Dictionary with serialized data. """ - return default_to_dict(self, document_store=self.document_store, cache_field=self.cache_field) + return default_to_dict( + self, + document_store=self.document_store, + cache_field=self.cache_field, + ttl=self.ttl.total_seconds() if self.ttl is not None else None, + time_field=self.time_field, + ) + + def _is_fresh(self, document: Document) -> bool: + """ + Checks whether a cached document is still within its TTL. + + Always returns True when no `ttl` is configured, to preserve the original non-expiring behavior. + + :param document: + The candidate cache-hit document. + :returns: + True if the document should count as a cache hit, False if it should be treated as expired. + """ + if self.ttl is None: + return True + + cached_at = document.meta.get(self.time_field) + if cached_at is None: + # ttl is enabled but the document was never stamped with a cache time: treat as stale + return False + + if isinstance(cached_at, str): + try: + cached_at = datetime.fromisoformat(cached_at) + except ValueError: + return False + + if not isinstance(cached_at, datetime): + return False + + if cached_at.tzinfo is None: + cached_at = cached_at.replace(tzinfo=timezone.utc) + + return datetime.now(timezone.utc) - cached_at < self.ttl @classmethod def from_dict(cls, data: dict[str, Any]) -> "CacheChecker": @@ -89,10 +138,12 @@ def run(self, items: list[Any]) -> dict[str, Any]: for item in items: filters = {"field": self.cache_field, "operator": "==", "value": item} found = self.document_store.filter_documents(filters=filters) - if found: - found_documents.extend(found) + fresh = [doc for doc in found if self._is_fresh(doc)] + if fresh: + found_documents.extend(fresh) else: misses.append(item) + return {"hits": found_documents, "misses": misses} @component.output_types(hits=list[Document], misses=list) @@ -116,8 +167,9 @@ async def run_async(self, items: list[Any]) -> dict[str, Any]: for item in items: filters = {"field": self.cache_field, "operator": "==", "value": item} found = await self.document_store.filter_documents_async(filters=filters) - if found: - found_documents.extend(found) + fresh = [doc for doc in found if self._is_fresh(doc)] + if fresh: + found_documents.extend(fresh) else: misses.append(item) return {"hits": found_documents, "misses": misses} diff --git a/releasenotes/notes/cache-checker-ttl-6842b76bcab8efd4.yaml b/releasenotes/notes/cache-checker-ttl-6842b76bcab8efd4.yaml new file mode 100644 index 00000000000..d487b2a39fd --- /dev/null +++ b/releasenotes/notes/cache-checker-ttl-6842b76bcab8efd4.yaml @@ -0,0 +1,5 @@ +--- +enhancements: + Add optional TTL support to ``CacheChecker`` to allow cached documents + to expire after a configurable amount of time. The cache timestamp + field can be customized using the ``time_field`` parameter. diff --git a/test/components/caching/test_cache_checker.py b/test/components/caching/test_cache_checker.py index 4ff3ff5bec7..12fa4e9e42c 100644 --- a/test/components/caching/test_cache_checker.py +++ b/test/components/caching/test_cache_checker.py @@ -2,6 +2,7 @@ # # SPDX-License-Identifier: Apache-2.0 +from datetime import datetime, timedelta, timezone from unittest.mock import Mock, patch import pytest @@ -22,21 +23,36 @@ def test_to_dict(self): "init_parameters": { "document_store": {"type": "haystack.testing.factory.MockedDocumentStore", "init_parameters": {}}, "cache_field": "url", + "ttl": None, + "time_field": "cached_at", }, } def test_to_dict_with_custom_init_parameters(self): mocked_docstore_class = document_store_class("MockedDocumentStore") - component = CacheChecker(document_store=mocked_docstore_class(), cache_field="my_url_field") + component = CacheChecker( + document_store=mocked_docstore_class(), + cache_field="my_url_field", + ttl=timedelta(hours=1), + time_field="my_time_field", + ) data = component.to_dict() assert data == { "type": "haystack.components.caching.cache_checker.CacheChecker", "init_parameters": { "document_store": {"type": "haystack.testing.factory.MockedDocumentStore", "init_parameters": {}}, "cache_field": "my_url_field", + "ttl": 3600.0, + "time_field": "my_time_field", }, } + def test_to_dict_with_numeric_ttl(self): + mocked_docstore_class = document_store_class("MockedDocumentStore") + component = CacheChecker(document_store=mocked_docstore_class(), cache_field="url", ttl=90) + data = component.to_dict() + assert data["init_parameters"]["ttl"] == 90.0 + def test_from_dict(self): data = { "type": "haystack.components.caching.cache_checker.CacheChecker", @@ -46,11 +62,15 @@ def test_from_dict(self): "init_parameters": {}, }, "cache_field": "my_url_field", + "ttl": 3600.0, + "time_field": "my_time_field", }, } component = CacheChecker.from_dict(data) assert isinstance(component.document_store, InMemoryDocumentStore) assert component.cache_field == "my_url_field" + assert component.ttl == timedelta(hours=1) + assert component.time_field == "my_time_field" def test_from_dict_without_docstore(self): data = {"type": "haystack.components.caching.cache_checker.CacheChecker", "init_parameters": {}} @@ -103,3 +123,84 @@ def test_close(self): checker = CacheChecker(document_store=nonclosable_document_store, cache_field="url") checker.close() assert nonclosable_document_store.mock_calls == [] + + def test_run_with_ttl_fresh_hit(self, in_memory_doc_store): + fresh_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now(timezone.utc) - timedelta(minutes=5)}, + ) + in_memory_doc_store.write_documents([fresh_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [fresh_doc], "misses": []} + + def test_run_with_ttl_expired_is_miss(self, in_memory_doc_store): + stale_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now(timezone.utc) - timedelta(hours=2)}, + ) + in_memory_doc_store.write_documents([stale_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [], "misses": ["https://example.com/1"]} + + def test_run_with_ttl_missing_time_field_is_miss(self, in_memory_doc_store): + undated_doc = Document(content="doc1", meta={"url": "https://example.com/1"}) + in_memory_doc_store.write_documents([undated_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [], "misses": ["https://example.com/1"]} + + def test_run_with_ttl_iso_string_timestamp(self, in_memory_doc_store): + fresh_doc = Document( + content="doc1", + meta={ + "url": "https://example.com/1", + "cached_at": (datetime.now(timezone.utc) - timedelta(minutes=5)).isoformat(), + }, + ) + in_memory_doc_store.write_documents([fresh_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [fresh_doc], "misses": []} + + def test_run_without_ttl_ignores_time_field(self, in_memory_doc_store): + # backward compatibility: no ttl configured means entries never expire, regardless of cached_at + old_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now(timezone.utc) - timedelta(days=365)}, + ) + in_memory_doc_store.write_documents([old_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url") + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [old_doc], "misses": []} + + def test_run_with_ttl_naive_datetime_timestamp(self, in_memory_doc_store): + # meta timestamp with no tzinfo at all (not datetime.now(timezone.utc)) + fresh_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now() - timedelta(minutes=5)}, # noqa: DTZ005 + ) + in_memory_doc_store.write_documents([fresh_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [fresh_doc], "misses": []} + + def test_run_with_ttl_malformed_iso_string_is_miss(self, in_memory_doc_store): + bad_doc = Document(content="doc1", meta={"url": "https://example.com/1", "cached_at": "not-a-timestamp"}) + in_memory_doc_store.write_documents([bad_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [], "misses": ["https://example.com/1"]} + + def test_run_with_ttl_non_datetime_timestamp_is_miss(self, in_memory_doc_store): + bad_doc = Document(content="doc1", meta={"url": "https://example.com/1", "cached_at": 12345}) + in_memory_doc_store.write_documents([bad_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = checker.run(items=["https://example.com/1"]) + assert results == {"hits": [], "misses": ["https://example.com/1"]} + + def test_ttl_accepts_numeric_seconds(self): + mocked_docstore_class = document_store_class("MockedDocumentStore") + checker = CacheChecker(document_store=mocked_docstore_class(), cache_field="url", ttl=3600) + assert checker.ttl == timedelta(hours=1) diff --git a/test/components/caching/test_cache_checker_async.py b/test/components/caching/test_cache_checker_async.py index d33307af875..8f956da79d1 100644 --- a/test/components/caching/test_cache_checker_async.py +++ b/test/components/caching/test_cache_checker_async.py @@ -2,6 +2,7 @@ # # SPDX-License-Identifier: Apache-2.0 +from datetime import datetime, timedelta, timezone from unittest.mock import AsyncMock, MagicMock, Mock import pytest @@ -73,3 +74,25 @@ async def test_close_async(self): checker = CacheChecker(document_store=nonclosable_document_store, cache_field="url") await checker.close_async() assert nonclosable_document_store.mock_calls == [] + + @pytest.mark.asyncio + async def test_run_async_with_ttl_fresh_hit(self, in_memory_doc_store): + fresh_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now(timezone.utc) - timedelta(minutes=5)}, + ) + in_memory_doc_store.write_documents([fresh_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = await checker.run_async(items=["https://example.com/1"]) + assert results == {"hits": [fresh_doc], "misses": []} + + @pytest.mark.asyncio + async def test_run_async_with_ttl_expired_is_miss(self, in_memory_doc_store): + stale_doc = Document( + content="doc1", + meta={"url": "https://example.com/1", "cached_at": datetime.now(timezone.utc) - timedelta(hours=2)}, + ) + in_memory_doc_store.write_documents([stale_doc]) + checker = CacheChecker(in_memory_doc_store, cache_field="url", ttl=timedelta(hours=1)) + results = await checker.run_async(items=["https://example.com/1"]) + assert results == {"hits": [], "misses": ["https://example.com/1"]}