feat(acms): implement Tantivy text search backend #1161
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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",
|
||||
]
|
||||
Reference in New Issue
Block a user