feat(acms): implement Tantivy text search backend #1161

Closed
aditya wants to merge 1 commits from feature/m6-tantivy-backend into master
8 changed files with 781 additions and 2 deletions
+6
View File
@@ -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,
+52
View File
@@ -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
+41
View File
@@ -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
+1
View File
@@ -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",
+41 -2
View File
@@ -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",
]