From 84fd4a1a2571cb1e1e9f426392a8479dea1b2b47 Mon Sep 17 00:00:00 2001 From: Aditya Chhabra Date: Mon, 30 Mar 2026 09:23:51 +0000 Subject: [PATCH] feat(acms): implement Tantivy text search backend Add a persistent Tantivy-backed ACMS text index/query implementation with real schema-backed indexing, term-filtered searches, and DI fallback to in-memory backends when Tantivy is unavailable. Address review feedback by removing the in-memory-only placeholder search path, reopening the on-disk index in Behave coverage, and keeping the issue #870 changelog and benchmark updates in the same atomic commit. ISSUES CLOSED: #870 --- CHANGELOG.md | 6 + benchmarks/bench_tantivy_text_backend.py | 52 +++ features/steps/tantivy_text_backend_steps.py | 191 ++++++++ features/tantivy_text_backend.feature | 41 ++ pyproject.toml | 1 + src/cleveragents/application/container.py | 43 +- .../infrastructure/acms/__init__.py | 13 + .../infrastructure/acms/tantivy_backends.py | 436 ++++++++++++++++++ 8 files changed, 781 insertions(+), 2 deletions(-) create mode 100644 benchmarks/bench_tantivy_text_backend.py create mode 100644 features/steps/tantivy_text_backend_steps.py create mode 100644 features/tantivy_text_backend.feature create mode 100644 src/cleveragents/infrastructure/acms/__init__.py create mode 100644 src/cleveragents/infrastructure/acms/tantivy_backends.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 36aab2df1..1861a839a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -268,6 +268,12 @@ non-existent ``container.resolve()``, matching corrected production code. Activated regression-guard BDD scenarios for ``plan tree``, ``plan explain``, and ``plan correct``. (#647) +- Implemented a Tantivy-backed ACMS text search backend with DI container + registration and graceful in-memory fallback when Tantivy is unavailable. + Added `TantivyTextIndexBackend` and `TantivyTextBackend` under + `cleveragents.infrastructure.acms`, wired text backend selection in the + application container, added Behave scenarios for index/search/remove/fallback, + and added an ASV benchmark suite for text-search latency. (#870) - Added the production ACMS skeleton compression stage via `DepthReductionCompressor`. The pipeline now re-renders inherited parent fragments to overview depths 0-1 using the UKO detail-level map chain, diff --git a/benchmarks/bench_tantivy_text_backend.py b/benchmarks/bench_tantivy_text_backend.py new file mode 100644 index 000000000..67121a520 --- /dev/null +++ b/benchmarks/bench_tantivy_text_backend.py @@ -0,0 +1,52 @@ +"""ASV benchmarks for Tantivy text backend search latency. + +Issue #870 requires a benchmark harness for text-search latency. +""" + +from __future__ import annotations + +import sys +from pathlib import Path +from tempfile import TemporaryDirectory + +_SRC = str(Path(__file__).resolve().parents[1] / "src") +if _SRC not in sys.path: + sys.path.insert(0, _SRC) + +from cleveragents.infrastructure.acms.tantivy_backends import ( # noqa: E402 + TantivyTextBackend, + TantivyTextIndexBackend, +) + + +class TantivyTextSearchSuite: + """Benchmark text search over a synthetic corpus.""" + + timeout = 120 + + def setup(self) -> None: + self._tmp = TemporaryDirectory(prefix="bench-tantivy-") + self._index_backend = TantivyTextIndexBackend( + index_dir=Path(self._tmp.name) / "index", + ) + self._query_backend = TantivyTextBackend(index_backend=self._index_backend) + for idx in range(2_000): + topic = "auth" if idx % 7 == 0 else "billing" + self._index_backend.index_document( + "local/bench", + f"uko://resource/{idx:05d}", + f"{topic} workflow handler number {idx}", + { + "resource_id": f"RES{idx:05d}", + "path": f"src/{topic}/{idx:05d}.py", + "uko_type": "code", + "language": "python", + }, + ) + + def time_search_auth_scope(self) -> None: + _ = self._query_backend.search( + "auth", + scope=frozenset({"RES00007", "RES00014", "RES00021"}), + max_results=20, + ) diff --git a/features/steps/tantivy_text_backend_steps.py b/features/steps/tantivy_text_backend_steps.py new file mode 100644 index 000000000..a465bce74 --- /dev/null +++ b/features/steps/tantivy_text_backend_steps.py @@ -0,0 +1,191 @@ +"""Step definitions for Tantivy text backend feature scenarios.""" + +from __future__ import annotations + +from pathlib import Path +from tempfile import TemporaryDirectory +from typing import Any + +from behave import given, then, when + +import cleveragents.application.container as container_module +from cleveragents.config.settings import get_settings +from cleveragents.domain.models.acms.index_stubs import InMemoryTextIndexBackend +from cleveragents.infrastructure.acms.tantivy_backends import ( + TantivyTextBackend, + TantivyTextIndexBackend, +) + + +def _new_backend_pair( + context: Any, +) -> tuple[TantivyTextIndexBackend, TantivyTextBackend]: + temp_dir = TemporaryDirectory(prefix="tantivy-backend-") + context._tantivy_tmp_dir = temp_dir + index_dir = Path(temp_dir.name) / "text-index" + index_backend = TantivyTextIndexBackend(index_dir=index_dir) + query_backend = TantivyTextBackend(index_backend=index_backend) + return index_backend, query_backend + + +def _index_scoped_document( + context: Any, + project: str, + doc_id: str, + resource_id: str, + path: str, +) -> None: + context.tantivy_index_backend.index_document( + project, + doc_id, + "handler function for request routing", + { + "resource_id": resource_id, + "path": path, + "uko_type": "code", + "language": "python", + }, + ) + + +@given("a fresh Tantivy text index backend") +def step_given_fresh_tantivy_text_index_backend(context: Any) -> None: + index_backend, _ = _new_backend_pair(context) + context.tantivy_index_backend = index_backend + + +@given("a fresh Tantivy text query backend") +def step_given_fresh_tantivy_text_query_backend(context: Any) -> None: + index_backend, query_backend = _new_backend_pair(context) + context.tantivy_index_backend = index_backend + context.tantivy_query_backend = query_backend + + +@when('I index project "{project}" document "{doc_id}" with content "{content}"') +def step_when_index_project_document( + context: Any, + project: str, + doc_id: str, + content: str, +) -> None: + context.tantivy_index_backend.index_document( + project, + doc_id, + content, + { + "resource_id": doc_id.rsplit("/", 1)[-1].upper(), + "path": "src/example.py", + "uko_type": "code", + "language": "python", + }, + ) + + +@when('I search indexed project "{project}" for "{query}"') +def step_when_search_indexed_project(context: Any, project: str, query: str) -> None: + context.index_search_results = context.tantivy_index_backend.search(project, query) + + +@then('indexed search should return document "{doc_id}"') +def step_then_indexed_search_returns_document(context: Any, doc_id: str) -> None: + ids = [result.doc_id for result in context.index_search_results] + assert doc_id in ids + + +@when( + 'I index scoped document "{doc_id}" with resource_id "{resource_id}" and path "{path}"' +) +def step_when_index_scoped_document( + context: Any, + doc_id: str, + resource_id: str, + path: str, +) -> None: + _index_scoped_document(context, "local/demo", doc_id, resource_id, path) + + +@when( + 'I index scoped project "{project}" document "{doc_id}" with resource_id "{resource_id}" and path "{path}"' +) +def step_when_index_scoped_project_document( + context: Any, + project: str, + doc_id: str, + resource_id: str, + path: str, +) -> None: + _index_scoped_document(context, project, doc_id, resource_id, path) + + +@when("I create a new Tantivy text query backend from the same index directory") +def step_when_create_new_query_backend_from_same_index(context: Any) -> None: + persisted_index_backend = TantivyTextIndexBackend( + index_dir=context.tantivy_index_backend.index_dir, + ) + context.tantivy_query_backend = TantivyTextBackend( + index_backend=persisted_index_backend, + ) + + +@when('I search text backend for "{query}" with scope "{scope_id}"') +def step_when_search_text_backend( + context: Any, + query: str, + scope_id: str, +) -> None: + context.text_search_results = context.tantivy_query_backend.search( + query, + scope=frozenset({scope_id}), + ) + + +@then('text search results should contain only "{doc_id}"') +def step_then_text_search_contains_only(context: Any, doc_id: str) -> None: + ids = [result.uko_uri for result in context.text_search_results] + assert ids == [doc_id] + + +@when('I remove project "{project}" document "{doc_id}"') +def step_when_remove_project_document(context: Any, project: str, doc_id: str) -> None: + context.tantivy_index_backend.remove_document(project, doc_id) + + +@then('index backend count for project "{project}" should be {count:d}') +def step_then_index_backend_count_for_project( + context: Any, + project: str, + count: int, +) -> None: + assert context.tantivy_index_backend.count(project) == count + + +@when('I clear index backend project "{project}"') +def step_when_clear_index_backend_project(context: Any, project: str) -> None: + context.tantivy_index_backend.clear(project) + + +@then("index backend total count should be {count:d}") +def step_then_index_backend_total_count(context: Any, count: int) -> None: + assert context.tantivy_index_backend.count() == count + + +@given("Tantivy availability is forced to false") +def step_given_tantivy_availability_forced_false(context: Any) -> None: + container_module_any: Any = container_module + context._orig_is_tantivy_available = container_module_any.is_tantivy_available + container_module_any.is_tantivy_available = lambda: False + + +@when("I build the text index backend from settings") +def step_when_build_text_index_backend_from_settings(context: Any) -> None: + settings = get_settings() + context.built_text_index_backend = container_module._build_text_index_backend( + settings, + ) + + +@then("the built index backend should be an InMemoryTextIndexBackend") +def step_then_built_backend_is_inmemory(context: Any) -> None: + assert isinstance(context.built_text_index_backend, InMemoryTextIndexBackend) + container_module_any: Any = container_module + container_module_any.is_tantivy_available = context._orig_is_tantivy_available diff --git a/features/tantivy_text_backend.feature b/features/tantivy_text_backend.feature new file mode 100644 index 000000000..7ef0ef903 --- /dev/null +++ b/features/tantivy_text_backend.feature @@ -0,0 +1,41 @@ +Feature: Tantivy text backend + As an ACMS developer + I want a production text backend for indexing and query + So that text-dependent strategies can retrieve relevant context + + Scenario: Index documents and query with Tantivy text index backend + Given a fresh Tantivy text index backend + When I index project "local/demo" document "uko://resource/alpha" with content "Auth flow handles login" + And I index project "local/demo" document "uko://resource/beta" with content "Billing flow handles invoices" + And I search indexed project "local/demo" for "auth" + Then indexed search should return document "uko://resource/alpha" + + Scenario: Query backend respects field filters and scope + Given a fresh Tantivy text query backend + When I index scoped document "uko://resource/scoped" with resource_id "RES01" and path "src/auth.py" + And I index scoped document "uko://resource/outside" with resource_id "RES02" and path "src/billing.py" + And I search text backend for "path:src/auth.py handler" with scope "RES01" + Then text search results should contain only "uko://resource/scoped" + + Scenario: Query backend reads documents from the persisted Tantivy index + Given a fresh Tantivy text index backend + When I index scoped project "local/demo" document "uko://resource/persisted" with resource_id "RES10" and path "src/acms.py" + And I create a new Tantivy text query backend from the same index directory + And I search text backend for "path:src/acms.py handler" with scope "RES10" + Then text search results should contain only "uko://resource/persisted" + + Scenario: Index backend supports remove and clear operations + Given a fresh Tantivy text index backend + When I index project "local/demo" document "uko://resource/remove-me" with content "temporary content" + And I remove project "local/demo" document "uko://resource/remove-me" + Then index backend count for project "local/demo" should be 0 + When I index project "local/demo" document "uko://resource/one" with content "alpha" + And I index project "local/other" document "uko://resource/two" with content "beta" + And I clear index backend project "local/demo" + Then index backend count for project "local/demo" should be 0 + And index backend total count should be 1 + + Scenario: Container builder degrades gracefully when Tantivy is unavailable + Given Tantivy availability is forced to false + When I build the text index backend from settings + Then the built index backend should be an InMemoryTextIndexBackend diff --git a/pyproject.toml b/pyproject.toml index ad642ab01..fad4def8d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -28,6 +28,7 @@ dependencies = [ "uvicorn>=0.30.1", "watchdog>=4.0.0", "faiss-cpu>=1.7.4", # Vector store backend + "tantivy>=0.22.0", # Text search backend "rx>=3.2.0", # Reactive streams for routing "dependency-injector>=4.41.0", # DI container "pydantic>=2.7.0", diff --git a/src/cleveragents/application/container.py b/src/cleveragents/application/container.py index 0c494505f..640b620ec 100644 --- a/src/cleveragents/application/container.py +++ b/src/cleveragents/application/container.py @@ -81,6 +81,8 @@ from cleveragents.application.services.uko_indexer import UKOIndexer from cleveragents.application.services.vector_store_service import VectorStoreService from cleveragents.config.settings import Settings, get_settings from cleveragents.domain.models.acms.analyzers import AnalyzerRegistry +from cleveragents.domain.models.acms.backends import TextBackend +from cleveragents.domain.models.acms.index_backends import TextIndexBackend from cleveragents.domain.models.acms.index_stubs import ( InMemoryGraphIndexBackend, InMemoryTextIndexBackend, @@ -90,6 +92,11 @@ from cleveragents.domain.models.acms.stubs import ( InMemoryTextBackend, ) from cleveragents.domain.providers.ai_provider import AIProviderInterface +from cleveragents.infrastructure.acms.tantivy_backends import ( + TantivyTextBackend, + TantivyTextIndexBackend, + is_tantivy_available, +) from cleveragents.infrastructure.database.llm_trace_repository import ( LLMTraceRepository, ) @@ -340,6 +347,32 @@ def _build_analyzer_registry() -> AnalyzerRegistry: return registry +def _build_text_index_backend(settings: Settings) -> TextIndexBackend: + """Build text index backend with Tantivy-first fallback.""" + text_index_dir = settings.data_dir / "index" / "text" + if is_tantivy_available(): + _logger.info( + "acms.text_index_backend.selected", + backend="tantivy", + index_dir=str(text_index_dir), + ) + return TantivyTextIndexBackend(index_dir=text_index_dir) + + _logger.warning( + "acms.text_index_backend.fallback", + backend="in_memory", + reason="tantivy_not_installed", + ) + return InMemoryTextIndexBackend() + + +def _build_text_query_backend(index_text_backend: TextIndexBackend) -> TextBackend: + """Build query-side text backend consistent with write backend.""" + if isinstance(index_text_backend, TantivyTextIndexBackend): + return TantivyTextBackend(index_backend=index_text_backend) + return InMemoryTextBackend() + + def _build_resource_file_watcher( event_bus: ReactiveEventBus, ) -> ResourceFileWatcher: @@ -742,7 +775,14 @@ class Container(containers.DeclarativeContainer): # ACMS Backend Abstraction Layer — configurable via provider selection. # Default: config-driven selection with graceful fallback to # in-memory stubs when production backends are disabled or missing. - text_backend = providers.Singleton(InMemoryTextBackend) + index_text_backend = providers.Singleton( + _build_text_index_backend, + settings=settings, + ) + text_backend = providers.Singleton( + _build_text_query_backend, + index_text_backend=index_text_backend, + ) vector_backend = providers.Singleton( build_vector_backend, vector_store_service=acms_vector_store_service, @@ -751,7 +791,6 @@ class Container(containers.DeclarativeContainer): # ACMS UKO Indexer — write-side index backends (#578) analyzer_registry = providers.Singleton(_build_analyzer_registry) - index_text_backend = providers.Singleton(InMemoryTextIndexBackend) index_vector_backend = providers.Singleton( build_vector_index_backend, vector_store_service=acms_vector_store_service, diff --git a/src/cleveragents/infrastructure/acms/__init__.py b/src/cleveragents/infrastructure/acms/__init__.py new file mode 100644 index 000000000..16a487c4d --- /dev/null +++ b/src/cleveragents/infrastructure/acms/__init__.py @@ -0,0 +1,13 @@ +"""ACMS infrastructure backends and adapters.""" + +from cleveragents.infrastructure.acms.tantivy_backends import ( + TantivyTextBackend, + TantivyTextIndexBackend, + is_tantivy_available, +) + +__all__: list[str] = [ + "TantivyTextBackend", + "TantivyTextIndexBackend", + "is_tantivy_available", +] diff --git a/src/cleveragents/infrastructure/acms/tantivy_backends.py b/src/cleveragents/infrastructure/acms/tantivy_backends.py new file mode 100644 index 000000000..2d0ba74a7 --- /dev/null +++ b/src/cleveragents/infrastructure/acms/tantivy_backends.py @@ -0,0 +1,436 @@ +"""Tantivy-backed ACMS text query and indexing backends.""" + +from __future__ import annotations + +from dataclasses import dataclass +from importlib import import_module +from pathlib import Path +from threading import RLock +from typing import Any + +import structlog + +from cleveragents.domain.models.acms.backends import TextResult +from cleveragents.domain.models.acms.index_backends import ( + IndexedDocument, + SearchResult, +) + +logger = structlog.get_logger(__name__) + +_QUERY_FIELDS: tuple[str, ...] = ("content",) +_METADATA_FIELDS: tuple[str, ...] = ( + "resource_id", + "path", + "uko_type", + "language", +) + + +def is_tantivy_available() -> bool: + """Return ``True`` when ``tantivy`` can be imported.""" + try: + import_module("tantivy") + except ModuleNotFoundError: + return False + return True + + +def _require_non_empty(value: str, name: str) -> None: + if not value or not value.strip(): + raise ValueError(f"{name} must be a non-empty string") + + +def _parse_field_filters(query: str) -> tuple[str, dict[str, str]]: + """Parse ``field:value`` filters and return free-text remainder.""" + filters: dict[str, str] = {} + free_terms: list[str] = [] + for token in query.split(): + if ":" in token: + field, value = token.split(":", 1) + key = field.strip().lower() + val = value.strip() + if key in {"path", "uko_type", "language", "uri"} and val: + filters[key] = val + continue + free_terms.append(token) + return " ".join(free_terms).strip(), filters + + +@dataclass(frozen=True) +class _IndexedTextDocument: + project: str + doc_id: str + content: str + metadata: dict[str, str] + + +class TantivyTextIndexBackend: + """Write-side text index backend backed by a persistent Tantivy index.""" + + def __init__( + self, + *, + index_dir: Path, + tantivy_module: Any | None = None, + ) -> None: + self._index_dir = index_dir.expanduser().resolve() + self._index_dir.mkdir(parents=True, exist_ok=True) + self._lock = RLock() + resolved_tantivy: Any = tantivy_module + if resolved_tantivy is None: + try: + resolved_tantivy = import_module("tantivy") + except ModuleNotFoundError as exc: + raise RuntimeError( + "tantivy is required for TantivyTextIndexBackend" + ) from exc + self._tantivy: Any = resolved_tantivy + self._logger = logger.bind(component="TantivyTextIndexBackend") + self._schema = self._build_schema() + self._index = self._tantivy.Index( + schema=self._schema, path=str(self._index_dir) + ) + + @property + def index_dir(self) -> Path: + return self._index_dir + + def index_document( + self, + project: str, + doc_id: str, + content: str, + metadata: dict[str, str], + ) -> IndexedDocument: + _require_non_empty(project, "project") + _require_non_empty(doc_id, "doc_id") + _require_non_empty(content, "content") + + record = _IndexedTextDocument( + project=project, + doc_id=doc_id, + content=content, + metadata=dict(metadata), + ) + with self._lock: + writer = self._index.writer(heap_size=50_000_000, num_threads=1) + self._delete_documents( + writer, + project=project, + filters={"uri": doc_id}, + ) + writer.add_document(self._build_document(record)) + self._commit(writer) + + return IndexedDocument(project=project, doc_id=doc_id, char_count=len(content)) + + def search( + self, + project: str, + query: str, + limit: int = 20, + ) -> list[SearchResult]: + _require_non_empty(project, "project") + _require_non_empty(query, "query") + if limit < 1: + raise ValueError(f"limit must be positive, got {limit}") + + free_text, filters = _parse_field_filters(query) + scored = self.search_documents( + free_text=free_text, + filters=filters, + limit=limit, + project=project, + ) + return [ + SearchResult( + doc_id=record.doc_id, + content=record.content[:200], + score=score, + metadata=dict(record.metadata), + ) + for score, record in scored[:limit] + ] + + def remove_document(self, project: str, doc_id: str) -> None: + _require_non_empty(project, "project") + _require_non_empty(doc_id, "doc_id") + with self._lock: + writer = self._index.writer(heap_size=50_000_000, num_threads=1) + self._delete_documents( + writer, + project=project, + filters={"uri": doc_id}, + ) + self._commit(writer) + + def rebuild_index(self, project: str) -> None: + _require_non_empty(project, "project") + self.clear(project) + + def clear(self, project: str | None = None) -> None: + """Clear all indexed docs, or just docs from one project.""" + with self._lock: + writer = self._index.writer(heap_size=50_000_000, num_threads=1) + if project is None: + writer.delete_all_documents() + self._commit(writer) + return + + _require_non_empty(project, "project") + self._delete_documents(writer, project=project) + self._commit(writer) + + def count(self, project: str | None = None) -> int: + return self.count_documents(project=project) + + def count_documents( + self, + *, + project: str | None = None, + filters: dict[str, str] | None = None, + scope: frozenset[str] | None = None, + ) -> int: + with self._lock: + query = self._build_query( + free_text="", + filters=filters or {}, + project=project, + scope=scope, + ) + self._index.reload() + searcher = self._index.searcher() + result = searcher.search(query, 1, count=True) + return int(result.count) + + def get_document(self, doc_id: str) -> _IndexedTextDocument | None: + _require_non_empty(doc_id, "doc_id") + with self._lock: + matches = self.search_documents( + free_text="", + filters={"uri": doc_id}, + limit=1, + ) + if not matches: + return None + return matches[0][1] + + def search_documents( + self, + *, + free_text: str, + filters: dict[str, str], + limit: int, + project: str | None = None, + scope: frozenset[str] | None = None, + ) -> list[tuple[float, _IndexedTextDocument]]: + if limit < 1: + raise ValueError(f"limit must be positive, got {limit}") + + with self._lock: + query = self._build_query( + free_text=free_text, + filters=filters, + project=project, + scope=scope, + ) + self._index.reload() + searcher = self._index.searcher() + result = searcher.search(query, limit) + + if not result.hits: + return [] + + max_score = max(float(score) for score, _ in result.hits) + normalized_max = max_score if max_score > 0.0 else 1.0 + documents: list[tuple[float, _IndexedTextDocument]] = [] + for raw_score, doc_address in result.hits: + stored = self._load_document(searcher.doc(doc_address)) + documents.append((min(float(raw_score) / normalized_max, 1.0), stored)) + return documents + + def _build_schema(self) -> Any: + builder = self._tantivy.SchemaBuilder() + for field_name in ( + "project", + "doc_id", + "resource_id", + "path", + "uko_type", + "language", + ): + builder.add_text_field(field_name, stored=True, tokenizer_name="raw") + builder.add_text_field("content", stored=True) + return builder.build() + + def _build_document(self, record: _IndexedTextDocument) -> Any: + document = self._tantivy.Document() + document.add_text("project", record.project) + document.add_text("doc_id", record.doc_id) + document.add_text("content", record.content) + for field_name in _METADATA_FIELDS: + value = record.metadata.get(field_name) + if value: + document.add_text(field_name, value) + return document + + def _build_query( + self, + *, + free_text: str, + filters: dict[str, str], + project: str | None = None, + scope: frozenset[str] | None = None, + ) -> Any: + clauses: list[Any] = [] + + if free_text: + clauses.append(self._index.parse_query(free_text, list(_QUERY_FIELDS))) + else: + clauses.append(self._tantivy.Query.all_query()) + + if project is not None: + _require_non_empty(project, "project") + clauses.append( + self._tantivy.Query.term_query(self._schema, "project", project) + ) + + if scope: + clauses.append( + self._tantivy.Query.term_set_query( + self._schema, + "resource_id", + sorted(scope), + ) + ) + + exact_filters = dict(filters) + if "uri" in exact_filters: + clauses.append( + self._tantivy.Query.term_query( + self._schema, + "doc_id", + exact_filters.pop("uri"), + ) + ) + for field_name, value in exact_filters.items(): + clauses.append( + self._tantivy.Query.term_query(self._schema, field_name, value) + ) + + if len(clauses) == 1: + return clauses[0] + + return self._tantivy.Query.boolean_query( + [(self._tantivy.Occur.Must, clause) for clause in clauses] + ) + + def _delete_documents( + self, + writer: Any, + *, + project: str | None = None, + filters: dict[str, str] | None = None, + ) -> None: + writer.delete_documents_by_query( + self._build_query( + free_text="", + filters=filters or {}, + project=project, + ) + ) + + def _commit(self, writer: Any) -> None: + writer.commit() + self._index.reload() + + def _load_document(self, document: Any) -> _IndexedTextDocument: + metadata = { + field_name: self._read_stored_text(document, field_name) + for field_name in _METADATA_FIELDS + if self._read_stored_text(document, field_name) + } + return _IndexedTextDocument( + project=self._read_stored_text(document, "project"), + doc_id=self._read_stored_text(document, "doc_id"), + content=self._read_stored_text(document, "content"), + metadata=metadata, + ) + + @staticmethod + def _read_stored_text(document: Any, field_name: str) -> str: + try: + values = document[field_name] + except KeyError: + return "" + if not values: + return "" + return str(values[0]) + + +class TantivyTextBackend: + """Read-side text backend backed by ``TantivyTextIndexBackend``.""" + + def __init__(self, *, index_backend: TantivyTextIndexBackend) -> None: + self._index_backend = index_backend + self._logger = logger.bind(component="TantivyTextBackend") + + def search( + self, + query: str, + *, + scope: frozenset[str], + max_results: int = 20, + ) -> list[TextResult]: + _require_non_empty(query, "query") + if max_results < 1: + raise ValueError(f"max_results must be positive, got {max_results}") + + free_text, filters = _parse_field_filters(query) + matches = self._index_backend.search_documents( + free_text=free_text, + filters=filters, + limit=max_results, + scope=scope, + ) + results = [ + TextResult( + uko_uri=doc.doc_id, + content=doc.content[:240], + score=score, + metadata=dict(doc.metadata), + ) + for score, doc in matches + ] + self._logger.debug( + "text.search.completed", + query=query, + result_count=len(results), + max_results=max_results, + ) + return results[:max_results] + + def get_by_uri(self, uko_uri: str) -> TextResult | None: + """Return one document by UKO URI/doc ID.""" + _require_non_empty(uko_uri, "uko_uri") + doc = self._index_backend.get_document(uko_uri) + if doc is None: + return None + return TextResult( + uko_uri=doc.doc_id, + content=doc.content, + score=1.0, + metadata=dict(doc.metadata), + ) + + def count(self, *, scope: frozenset[str] | None = None) -> int: + """Return indexed document count, optionally constrained to scope.""" + return self._index_backend.count_documents(scope=scope) + + +__all__: list[str] = [ + "TantivyTextBackend", + "TantivyTextIndexBackend", + "is_tantivy_available", +] -- 2.52.0