diff --git a/benchmarks/faiss_vector_backend_bench.py b/benchmarks/faiss_vector_backend_bench.py new file mode 100644 index 000000000..d62517eab --- /dev/null +++ b/benchmarks/faiss_vector_backend_bench.py @@ -0,0 +1,102 @@ +"""ASV benchmarks for the FAISS ACMS vector backend adapters.""" + +from __future__ import annotations + +import importlib +import os +import shutil +import sys +import tempfile +from pathlib import Path +from typing import ClassVar + +_SRC = str(Path(__file__).resolve().parents[1] / "src") +if _SRC not in sys.path: + sys.path.insert(0, _SRC) + +import cleveragents # noqa: E402 + +importlib.reload(cleveragents) + +from cleveragents.application.services.config_service import ConfigService # noqa: E402 +from cleveragents.application.services.faiss_vector_backend import ( # noqa: E402 + FAISSVectorBackend, + FAISSVectorIndexBackend, +) +from cleveragents.application.services.vector_store_service import ( # noqa: E402 + FAISS, + VectorStoreService, +) +from cleveragents.config.settings import Settings # noqa: E402 +from cleveragents.infrastructure.database.unit_of_work import UnitOfWork # noqa: E402 + + +class VectorSearchLatencySuite: + """Benchmark FAISS-backed ACMS vector indexing and search.""" + + params: ClassVar[list[int]] = [250, 1000] + param_names: ClassVar[list[str]] = ["documents"] + timeout = 120 + + def setup(self, documents: int) -> None: + if FAISS is None: + raise RuntimeError("FAISS is required for this benchmark") + + self._saved_env = { + key: os.environ.get(key) + for key in ( + "CLEVERAGENTS_INDEX_VECTOR_BACKEND", + "CLEVERAGENTS_INDEX_VECTOR_DIR", + "CLEVERAGENTS_EMBEDDING_PROVIDER", + "CLEVERAGENTS_EMBEDDING_DIMENSIONS", + ) + } + self._tempdir = tempfile.mkdtemp(prefix="asv-faiss-vector-") + os.environ["CLEVERAGENTS_INDEX_VECTOR_BACKEND"] = "faiss" + os.environ["CLEVERAGENTS_INDEX_VECTOR_DIR"] = self._tempdir + os.environ["CLEVERAGENTS_EMBEDDING_PROVIDER"] = "fake" + os.environ["CLEVERAGENTS_EMBEDDING_DIMENSIONS"] = "32" + + service = VectorStoreService( + Settings(), + UnitOfWork("sqlite:///:memory:", require_confirmation=False), + ConfigService(), + ) + self._index_backend = FAISSVectorIndexBackend(service) + self._vector_backend = FAISSVectorBackend(service) + self._query = [1.0, 2.0, 3.0] + self._scope = frozenset({f"RES{i:04d}" for i in range(min(documents, 32))}) + + for i in range(documents): + vector = [float(i % 11), float((i * 3) % 17), float((i * 5) % 19)] + self._index_backend.index_embedding( + "local/bench", + f"uko:{i:04d}", + vector, + { + "resource_id": f"RES{i:04d}", + "location": f"src/file_{i:04d}.py", + "resource_type": "python", + }, + ) + + def teardown(self, documents: int) -> None: + _ = documents + shutil.rmtree(self._tempdir, ignore_errors=True) + for key, value in self._saved_env.items(): + if value is None: + os.environ.pop(key, None) + else: + os.environ[key] = value + + def time_project_search(self, documents: int) -> None: + _ = documents + self._index_backend.search_similar("local/bench", self._query, limit=20) + + def time_scoped_similarity_search(self, documents: int) -> None: + _ = documents + self._vector_backend.similarity_search( + self._query, + scope=self._scope, + top_k=20, + ) diff --git a/features/faiss_vector_backend.feature b/features/faiss_vector_backend.feature new file mode 100644 index 000000000..80ed41add --- /dev/null +++ b/features/faiss_vector_backend.feature @@ -0,0 +1,46 @@ +Feature: FAISS ACMS vector backend + As an ACMS maintainer + I want the FAISS vector backend wired into the ACMS protocols + So that semantic vector search can index and retrieve UKO resources + + Scenario: FAISS vector backends index and search embeddings with scope filtering + Given the ACMS FAISS backend is configured to use fake embeddings + And FAISS interactions are recorded for ACMS backends + And an ACMS FAISS backend pair + When I index ACMS embedding "uko:auth" in project "local/app" for resource "RES01" with vector "1.0,0.0" + And I index ACMS embedding "uko:billing" in project "local/app" for resource "RES02" with vector "4.0,0.0" + And I search ACMS vectors in project "local/app" for vector "1.1,0.0" + Then the project vector search should include doc_id "uko:auth" + When I search ACMS vectors scoped to resource "RES01" for vector "1.1,0.0" + Then the scoped vector search should include uko_uri "uko:auth" + When I search ACMS vectors scoped to resource "RES02" for vector "1.1,0.0" + Then the scoped vector search should include uko_uri "uko:billing" + + Scenario: FAISS vector index backend removes embeddings + Given the ACMS FAISS backend is configured to use fake embeddings + And FAISS interactions are recorded for ACMS backends + And an ACMS FAISS backend pair + And I index ACMS embedding "uko:auth" in project "local/app" for resource "RES01" with vector "1.0,0.0" + When I remove ACMS embedding "uko:auth" from project "local/app" + Then project "local/app" should have 0 ACMS embeddings + + Scenario: ACMS FAISS backend uses embedding configuration from ConfigService + Given the ACMS FAISS backend is configured to use fake embeddings + And FAISS interactions are recorded for ACMS backends + And an ACMS FAISS backend pair + When I index ACMS embedding "uko:auth" in project "local/app" for resource "RES01" with vector "1.0,0.0" + Then the ACMS embeddings provider should be FakeEmbeddings sized 8 + + Scenario: DI container resolves FAISS vector backends when configured + Given the ACMS FAISS backend is configured to use fake embeddings + And FAISS interactions are recorded for ACMS backends + When I resolve ACMS vector backends from the DI container + Then the container vector backend should be a FAISSVectorBackend + And the container index vector backend should be a FAISSVectorIndexBackend + + Scenario: DI container falls back when FAISS support is unavailable + Given the ACMS FAISS backend is configured to use fake embeddings + And FAISS is unavailable for ACMS backends + When I resolve ACMS vector backends from the DI container + Then the container vector backend should be an InMemoryVectorBackend + And the container index vector backend should be an InMemoryVectorIndexBackend diff --git a/features/steps/faiss_vector_backend_steps.py b/features/steps/faiss_vector_backend_steps.py new file mode 100644 index 000000000..2d140df57 --- /dev/null +++ b/features/steps/faiss_vector_backend_steps.py @@ -0,0 +1,338 @@ +"""Step definitions for the FAISS ACMS vector backend feature.""" + +from __future__ import annotations + +import os +import shutil +import tempfile +from pathlib import Path +from typing import Any, ClassVar, cast +from unittest.mock import patch + +from behave import given, then, when +from langchain_community.embeddings import FakeEmbeddings + +from cleveragents.application.container import get_container, reset_container +from cleveragents.application.services.config_service import ConfigService +from cleveragents.application.services.faiss_vector_backend import ( + FAISSVectorBackend, + FAISSVectorIndexBackend, +) +from cleveragents.application.services.vector_store_service import VectorStoreService +from cleveragents.config.settings import Settings +from cleveragents.domain.models.acms.index_stubs import InMemoryVectorIndexBackend +from cleveragents.domain.models.acms.stubs import InMemoryVectorBackend +from cleveragents.infrastructure.database.unit_of_work import UnitOfWork +from features.steps.service_steps import add_cleanup + + +class _StubDocument: + def __init__(self, page_content: str, metadata: dict[str, str]) -> None: + self.page_content = page_content + self.metadata = dict(metadata) + + +class _StubDocStore: + def __init__(self) -> None: + self._docs: dict[str, _StubDocument] = {} + + def search(self, store_id: str) -> _StubDocument: + return self._docs[store_id] + + +class RecordingFAISS: + """Deterministic FAISS stand-in for ACMS backend tests.""" + + saved_snapshots: ClassVar[ + dict[str, list[tuple[str, str, list[float], dict[str, str]]]] + ] = {} + last_embedding_backend: ClassVar[Any | None] = None + + def __init__(self) -> None: + self.docstore = _StubDocStore() + self.index_to_docstore_id: dict[int, str] = {} + self._vectors: dict[str, list[float]] = {} + + @classmethod + def reset(cls) -> None: + cls.saved_snapshots = {} + cls.last_embedding_backend = None + + @classmethod + def from_embeddings( + cls, + text_embeddings: list[tuple[str, list[float]]], + embedding: Any, + metadatas: list[dict[str, str]] | None = None, + ids: list[str] | None = None, + **_: Any, + ) -> RecordingFAISS: + instance = cls() + cls.last_embedding_backend = embedding + instance.add_embeddings(text_embeddings, metadatas=metadatas, ids=ids) + return instance + + def add_embeddings( + self, + text_embeddings: list[tuple[str, list[float]]], + metadatas: list[dict[str, str]] | None = None, + ids: list[str] | None = None, + **_: Any, + ) -> list[str]: + stored_ids: list[str] = [] + metadata_list = metadatas or [{} for _ in text_embeddings] + for index, (text, vector) in enumerate(text_embeddings): + store_id = ids[index] if ids is not None else f"doc-{len(self._vectors)}" + self.docstore._docs[store_id] = _StubDocument(text, metadata_list[index]) + self._vectors[store_id] = list(vector) + stored_ids.append(store_id) + self._rebuild_index_map() + return stored_ids + + def similarity_search_with_score_by_vector( + self, + embedding: list[float], + k: int = 4, + filter: Any | None = None, + fetch_k: int = 20, + **_: Any, + ) -> list[tuple[_StubDocument, float]]: + matches: list[tuple[_StubDocument, float]] = [] + for store_id, vector in self._vectors.items(): + document = self.docstore.search(store_id) + if filter is not None and not filter(document.metadata): + continue + distance = sum( + abs(left - right) + for left, right in zip(embedding, vector, strict=False) + ) + matches.append((document, float(distance))) + matches.sort(key=lambda item: item[1]) + return matches[: min(fetch_k, k)] + + def save_local(self, directory: str, index_name: str = "index") -> None: + _ = index_name + RecordingFAISS.saved_snapshots[directory] = [ + ( + store_id, + self.docstore.search(store_id).page_content, + list(self._vectors[store_id]), + dict(self.docstore.search(store_id).metadata), + ) + for store_id in self.index_to_docstore_id.values() + ] + target = Path(directory) + target.mkdir(parents=True, exist_ok=True) + (target / "index.faiss").write_bytes(b"stub") + (target / "index.pkl").write_bytes(b"stub") + + @classmethod + def load_local( + cls, + directory: str, + embeddings: Any, + index_name: str = "index", + *, + allow_dangerous_deserialization: bool = False, + **_: Any, + ) -> RecordingFAISS: + _unused = (index_name, allow_dangerous_deserialization) + cls.last_embedding_backend = embeddings + if directory not in cls.saved_snapshots: + raise FileNotFoundError(directory) + instance = cls() + for store_id, text, vector, metadata in cls.saved_snapshots[directory]: + instance.docstore._docs[store_id] = _StubDocument(text, metadata) + instance._vectors[store_id] = list(vector) + instance._rebuild_index_map() + return instance + + def delete(self, ids: list[str] | None = None, **_: Any) -> bool: + for store_id in ids or []: + self.docstore._docs.pop(store_id, None) + self._vectors.pop(store_id, None) + self._rebuild_index_map() + return True + + def _rebuild_index_map(self) -> None: + ordered_ids = sorted(self._vectors) + self.index_to_docstore_id = { + position: store_id for position, store_id in enumerate(ordered_ids) + } + + +def _set_env(context: Any, key: str, value: str | None) -> None: + original = os.environ.get(key) + if value is None: + os.environ.pop(key, None) + else: + os.environ[key] = value + + def cleanup() -> None: + if original is None: + os.environ.pop(key, None) + else: + os.environ[key] = original + + add_cleanup(context, cleanup) + + +def _parse_vector(raw: str) -> list[float]: + return [float(part.strip()) for part in raw.split(",") if part.strip()] + + +@given("the ACMS FAISS backend is configured to use fake embeddings") +def step_configure_acms_faiss(context: Any) -> None: + index_dir = tempfile.mkdtemp(prefix="acms-faiss-index-") + add_cleanup(context, lambda: shutil.rmtree(index_dir, ignore_errors=True)) + _set_env(context, "CLEVERAGENTS_INDEX_VECTOR_BACKEND", "faiss") + _set_env(context, "CLEVERAGENTS_INDEX_VECTOR_DIR", index_dir) + _set_env(context, "CLEVERAGENTS_EMBEDDING_PROVIDER", "fake") + _set_env(context, "CLEVERAGENTS_EMBEDDING_DIMENSIONS", "8") + reset_container() + + +@given("FAISS interactions are recorded for ACMS backends") +def step_patch_acms_faiss(context: Any) -> None: + RecordingFAISS.reset() + patchers = [ + patch( + "cleveragents.application.services.vector_store_service.FAISS", + RecordingFAISS, + ), + patch( + "cleveragents.application.services.faiss_vector_backend.FAISS", + RecordingFAISS, + ), + ] + for patcher in patchers: + patcher.start() + add_cleanup(context, patcher.stop) + + +@given("FAISS is unavailable for ACMS backends") +def step_disable_acms_faiss(context: Any) -> None: + patchers = [ + patch("cleveragents.application.services.vector_store_service.FAISS", None), + patch("cleveragents.application.services.faiss_vector_backend.FAISS", None), + ] + for patcher in patchers: + patcher.start() + add_cleanup(context, patcher.stop) + reset_container() + + +@given("an ACMS FAISS backend pair") +def step_acms_backend_pair(context: Any) -> None: + settings = Settings() + unit_of_work = UnitOfWork("sqlite:///:memory:", require_confirmation=False) + service = VectorStoreService(settings, unit_of_work, ConfigService()) + context.acms_service = service + context.acms_vector_backend = FAISSVectorBackend(service) + context.acms_index_backend = FAISSVectorIndexBackend(service) + + +@given( + 'I index ACMS embedding "{doc_id}" in project "{project}" for resource "{resource_id}" with vector "{vector}"' +) +@when( + 'I index ACMS embedding "{doc_id}" in project "{project}" for resource "{resource_id}" with vector "{vector}"' +) +def step_index_acms_embedding( + context: Any, + doc_id: str, + project: str, + resource_id: str, + vector: str, +) -> None: + context.acms_index_backend.index_embedding( + project, + doc_id, + _parse_vector(vector), + { + "resource_id": resource_id, + "location": f"src/{resource_id.lower()}.py", + "resource_type": "python", + }, + ) + + +@when('I search ACMS vectors in project "{project}" for vector "{vector}"') +def step_search_project_vectors(context: Any, project: str, vector: str) -> None: + context.index_results = context.acms_index_backend.search_similar( + project, + _parse_vector(vector), + limit=5, + ) + + +@when('I search ACMS vectors scoped to resource "{resource_id}" for vector "{vector}"') +def step_search_scoped_vectors(context: Any, resource_id: str, vector: str) -> None: + context.vector_results = context.acms_vector_backend.similarity_search( + _parse_vector(vector), + scope=frozenset({resource_id}), + top_k=5, + ) + + +@when('I remove ACMS embedding "{doc_id}" from project "{project}"') +def step_remove_acms_embedding(context: Any, doc_id: str, project: str) -> None: + context.acms_index_backend.remove_embedding(project, doc_id) + + +@when("I resolve ACMS vector backends from the DI container") +def step_resolve_container_backends(context: Any) -> None: + reset_container() + container = get_container() + context.container_vector_backend = container.vector_backend() + context.container_index_vector_backend = container.index_vector_backend() + + +@then('the project vector search should include doc_id "{doc_id}"') +def step_assert_project_result(context: Any, doc_id: str) -> None: + assert any(result.doc_id == doc_id for result in context.index_results) + + +@then('the scoped vector search should include uko_uri "{uko_uri}"') +def step_assert_scoped_result(context: Any, uko_uri: str) -> None: + assert any(result.uko_uri == uko_uri for result in context.vector_results) + + +@then("the scoped vector search should be empty") +def step_assert_scoped_empty(context: Any) -> None: + assert context.vector_results == [] + + +@then('project "{project}" should have {count:d} ACMS embeddings') +def step_assert_embedding_count(context: Any, project: str, count: int) -> None: + assert context.acms_service.acms_count(project=project) == count + + +@then("the ACMS embeddings provider should be FakeEmbeddings sized {size:d}") +def step_assert_fake_embeddings(context: Any, size: int) -> None: + backend = RecordingFAISS.last_embedding_backend + assert isinstance(backend, FakeEmbeddings) + typed_backend = cast(Any, backend) + assert typed_backend.size == size + + +@then("the container vector backend should be a FAISSVectorBackend") +def step_assert_container_vector_backend(context: Any) -> None: + assert isinstance(context.container_vector_backend, FAISSVectorBackend) + + +@then("the container index vector backend should be a FAISSVectorIndexBackend") +def step_assert_container_index_backend(context: Any) -> None: + assert isinstance(context.container_index_vector_backend, FAISSVectorIndexBackend) + + +@then("the container vector backend should be an InMemoryVectorBackend") +def step_assert_container_vector_fallback(context: Any) -> None: + assert isinstance(context.container_vector_backend, InMemoryVectorBackend) + + +@then("the container index vector backend should be an InMemoryVectorIndexBackend") +def step_assert_container_index_fallback(context: Any) -> None: + assert isinstance( + context.container_index_vector_backend, InMemoryVectorIndexBackend + ) diff --git a/src/cleveragents/application/container.py b/src/cleveragents/application/container.py index 094c216a7..a5c099470 100644 --- a/src/cleveragents/application/container.py +++ b/src/cleveragents/application/container.py @@ -30,6 +30,7 @@ from cleveragents.application.services.autonomy_guardrail_service import ( AutonomyGuardrailService, ) from cleveragents.application.services.checkpoint_service import CheckpointService +from cleveragents.application.services.config_service import ConfigService from cleveragents.application.services.context_service import ContextService from cleveragents.application.services.context_tiers import ( ContextTierService, @@ -43,6 +44,10 @@ from cleveragents.application.services.decomposition_service import ( from cleveragents.application.services.execution_environment_resolver import ( ExecutionEnvironmentResolver, ) +from cleveragents.application.services.faiss_vector_backend import ( + build_vector_backend, + build_vector_index_backend, +) from cleveragents.application.services.fix_then_revalidate import ( FixThenRevalidateOrchestrator, ) @@ -79,12 +84,10 @@ from cleveragents.domain.models.acms.analyzers import AnalyzerRegistry from cleveragents.domain.models.acms.index_stubs import ( InMemoryGraphIndexBackend, InMemoryTextIndexBackend, - InMemoryVectorIndexBackend, ) from cleveragents.domain.models.acms.stubs import ( InMemoryGraphBackend, InMemoryTextBackend, - InMemoryVectorBackend, ) from cleveragents.domain.providers.ai_provider import AIProviderInterface from cleveragents.infrastructure.database.llm_trace_repository import ( @@ -519,6 +522,7 @@ class Container(containers.DeclarativeContainer): UnitOfWork, database_url=database_url, ) + config_service = providers.Singleton(ConfigService) # AI Provider - Singleton that builds either a mock or real provider ai_provider: providers.Provider[AIProviderInterface | None] = providers.Singleton( @@ -542,6 +546,13 @@ class Container(containers.DeclarativeContainer): VectorStoreService, settings=settings, unit_of_work=unit_of_work, + config_service=config_service, + ) + acms_vector_store_service = providers.Singleton( + VectorStoreService, + settings=settings, + unit_of_work=unit_of_work, + config_service=config_service, ) context_service = providers.Factory( @@ -729,16 +740,22 @@ class Container(containers.DeclarativeContainer): plugin_manager = providers.Singleton(PluginManager) # ACMS Backend Abstraction Layer — configurable via provider selection. - # Default: in-memory stubs. Override with production backends - # (Tantivy, FAISS, Blazegraph, etc.) via ``override_providers()``. + # Default: config-driven selection with graceful fallback to + # in-memory stubs when production backends are disabled or missing. text_backend = providers.Singleton(InMemoryTextBackend) - vector_backend = providers.Singleton(InMemoryVectorBackend) + vector_backend = providers.Singleton( + build_vector_backend, + vector_store_service=acms_vector_store_service, + ) graph_backend = providers.Singleton(InMemoryGraphBackend) # 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(InMemoryVectorIndexBackend) + index_vector_backend = providers.Singleton( + build_vector_index_backend, + vector_store_service=acms_vector_store_service, + ) index_graph_backend = providers.Singleton(InMemoryGraphIndexBackend) uko_indexer = providers.Singleton( UKOIndexer, diff --git a/src/cleveragents/application/services/faiss_vector_backend.py b/src/cleveragents/application/services/faiss_vector_backend.py new file mode 100644 index 000000000..71175db15 --- /dev/null +++ b/src/cleveragents/application/services/faiss_vector_backend.py @@ -0,0 +1,204 @@ +"""FAISS-backed ACMS vector backend adapters. + +Bridges the ACMS read/write vector backend protocols to the shared +``VectorStoreService`` FAISS persistence helpers. +""" + +from __future__ import annotations + +from typing import Any + +import structlog + +from cleveragents.application.services.vector_store_service import ( + FAISS, + VectorStoreService, +) +from cleveragents.domain.models.acms.backends import VectorResult +from cleveragents.domain.models.acms.index_backends import SearchResult +from cleveragents.domain.models.acms.index_stubs import InMemoryVectorIndexBackend +from cleveragents.domain.models.acms.stubs import InMemoryVectorBackend + +logger = structlog.get_logger(__name__) + + +def _require_non_empty(value: str, name: str) -> str: + trimmed = value.strip() + if not trimmed: + raise ValueError(f"{name} must be a non-empty string") + return trimmed + + +class FAISSVectorBackend: + """Read-side ACMS vector backend backed by the shared FAISS store.""" + + def __init__(self, vector_store_service: VectorStoreService) -> None: + self._vector_store_service = vector_store_service + + def similarity_search( + self, + embedding: list[float], + *, + scope: frozenset[str], + top_k: int = 20, + ) -> list[VectorResult]: + """Search the shared ACMS FAISS store.""" + + if not embedding: + raise ValueError("embedding must be a non-empty list") + if top_k < 1: + raise ValueError(f"top_k must be positive, got {top_k}") + + hits = self._vector_store_service.acms_search_by_vector( + list(embedding), + top_k=top_k, + scope=scope, + ) + results: list[VectorResult] = [] + for hit in hits: + metadata = dict(hit["metadata"]) + doc_id = metadata.pop("doc_id", "") + if not doc_id: + continue + results.append( + VectorResult( + uko_uri=doc_id, + content=str(hit["content"]), + score=float(hit["score"]), + metadata=metadata, + ) + ) + return results + + +class FAISSVectorIndexBackend: + """Write-side ACMS vector index backend backed by the shared FAISS store.""" + + def __init__(self, vector_store_service: VectorStoreService) -> None: + self._vector_store_service = vector_store_service + + def index_embedding( + self, + project: str, + doc_id: str, + embedding: list[float], + metadata: dict[str, str], + ) -> None: + """Index an embedding in the shared ACMS FAISS store.""" + + normalized_project = _require_non_empty(project, "project") + normalized_doc_id = _require_non_empty(doc_id, "doc_id") + if not embedding: + raise ValueError("embedding must be a non-empty list") + + payload = metadata.get("location") or metadata.get( + "resource_type", f"embedding:{normalized_doc_id}" + ) + self._vector_store_service.acms_index_embedding( + normalized_project, + normalized_doc_id, + list(embedding), + dict(metadata), + content=payload, + ) + + def search_similar( + self, + project: str, + query_embedding: list[float], + limit: int = 20, + min_relevance: float = 0.0, + ) -> list[SearchResult]: + """Search embeddings within one project.""" + + normalized_project = _require_non_empty(project, "project") + if not query_embedding: + raise ValueError("query_embedding must be a non-empty list") + if limit < 1: + raise ValueError(f"limit must be positive, got {limit}") + if not (0.0 <= min_relevance <= 1.0): + raise ValueError( + f"min_relevance must be in [0.0, 1.0], got {min_relevance}" + ) + + hits = self._vector_store_service.acms_search_by_vector( + list(query_embedding), + top_k=limit, + project=normalized_project, + ) + results: list[SearchResult] = [] + for hit in hits: + score = float(hit["score"]) + if score < min_relevance: + continue + metadata = dict(hit["metadata"]) + doc_id = metadata.pop("doc_id", "") + if not doc_id: + continue + results.append( + SearchResult( + doc_id=doc_id, + content=str(hit["content"]), + score=score, + metadata=metadata, + ) + ) + return results + + def remove_embedding(self, project: str, doc_id: str) -> None: + """Remove one embedding from the shared ACMS FAISS store.""" + + normalized_project = _require_non_empty(project, "project") + normalized_doc_id = _require_non_empty(doc_id, "doc_id") + self._vector_store_service.acms_remove_embedding( + normalized_project, + normalized_doc_id, + ) + + +def build_vector_backend(vector_store_service: VectorStoreService) -> Any: + """Return the configured ACMS read-side vector backend.""" + + backend_name = vector_store_service.acms_backend_name() + if backend_name != "faiss": + logger.info( + "acms.vector_backend.fallback", + configured_backend=backend_name, + fallback_backend="InMemoryVectorBackend", + ) + return InMemoryVectorBackend() + if FAISS is None: + logger.warning( + "acms.vector_backend.faiss_unavailable", + fallback_backend="InMemoryVectorBackend", + ) + return InMemoryVectorBackend() + return FAISSVectorBackend(vector_store_service) + + +def build_vector_index_backend(vector_store_service: VectorStoreService) -> Any: + """Return the configured ACMS write-side vector backend.""" + + backend_name = vector_store_service.acms_backend_name() + if backend_name != "faiss": + logger.info( + "acms.vector_index_backend.fallback", + configured_backend=backend_name, + fallback_backend="InMemoryVectorIndexBackend", + ) + return InMemoryVectorIndexBackend() + if FAISS is None: + logger.warning( + "acms.vector_index_backend.faiss_unavailable", + fallback_backend="InMemoryVectorIndexBackend", + ) + return InMemoryVectorIndexBackend() + return FAISSVectorIndexBackend(vector_store_service) + + +__all__: list[str] = [ + "FAISSVectorBackend", + "FAISSVectorIndexBackend", + "build_vector_backend", + "build_vector_index_backend", +] diff --git a/src/cleveragents/application/services/vector_store_service.py b/src/cleveragents/application/services/vector_store_service.py index 9fdbe67e7..b24fae2a3 100644 --- a/src/cleveragents/application/services/vector_store_service.py +++ b/src/cleveragents/application/services/vector_store_service.py @@ -1,30 +1,103 @@ -"""Vector store management for semantic context search.""" +"""Vector store management for semantic context search. + +Supports two FAISS-backed use cases: + +- plan-scoped semantic search used by :class:`ContextService` +- ACMS vector indexing/query adapters used by the backend protocols +""" from __future__ import annotations from collections.abc import Iterable from pathlib import Path -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, Protocol, cast from langchain_community.embeddings import FakeEmbeddings -from langchain_community.vectorstores.faiss import FAISS + +try: + from langchain_community.vectorstores.faiss import FAISS +except ImportError: # pragma: no cover - exercised via graceful fallback tests + FAISS = None if TYPE_CHECKING: from langchain_core.embeddings import Embeddings +from cleveragents.application.services.config_service import ConfigService from cleveragents.config.settings import Settings from cleveragents.core.exceptions import ConfigurationError from cleveragents.domain.models.core import Context from cleveragents.infrastructure.database.unit_of_work import UnitOfWork +class _FaissPlanStoreProtocol(Protocol): + def similarity_search_with_score( + self, + query: str, + *, + k: int, + ) -> list[tuple[Any, Any]]: + _ = (query, k) + raise NotImplementedError + + def save_local(self, folder_path: str) -> None: + _ = folder_path + raise NotImplementedError + + +class _FaissDocstoreProtocol(Protocol): + def search(self, key: str) -> Any: + _ = key + raise NotImplementedError + + +class _FaissAcmsStoreProtocol(_FaissPlanStoreProtocol, Protocol): + index_to_docstore_id: dict[int, str] + docstore: _FaissDocstoreProtocol + + def add_embeddings( + self, + text_embeddings: list[tuple[str, list[float]]], + *, + metadatas: list[dict[str, str]], + ids: list[str], + ) -> Any: + _ = (text_embeddings, metadatas, ids) + raise NotImplementedError + + def similarity_search_with_score_by_vector( + self, + embedding: list[float], + *, + k: int, + filter: Any = None, + fetch_k: int, + ) -> list[tuple[Any, Any]]: + _ = (embedding, k, filter, fetch_k) + raise NotImplementedError + + def delete(self, *, ids: list[str]) -> Any: + _ = ids + raise NotImplementedError + + def save_local(self, folder_path: str) -> None: + _ = folder_path + raise NotImplementedError + + class VectorStoreService: """Manage creation and querying of LangChain vector stores.""" - def __init__(self, settings: Settings, unit_of_work: UnitOfWork): + def __init__( + self, + settings: Settings, + unit_of_work: UnitOfWork, + config_service: ConfigService | None = None, + ) -> None: self.settings = settings self.unit_of_work = unit_of_work - self._cache: dict[int, FAISS] = {} + self._config_service = config_service or ConfigService() + self._cache: dict[int, _FaissPlanStoreProtocol] = {} + self._acms_store: _FaissAcmsStoreProtocol | None = None # ------------------------------------------------------------------ # Public API @@ -55,7 +128,8 @@ class VectorStoreService: return 0 embeddings = self._create_embeddings() - vector_store = FAISS.from_texts( + faiss_cls = self._require_faiss() + vector_store = faiss_cls.from_texts( documents, embedding=embeddings, metadatas=metadata, @@ -103,6 +177,162 @@ class VectorStoreService: ) return formatted + # ------------------------------------------------------------------ + # ACMS vector index API + # ------------------------------------------------------------------ + def acms_backend_name(self) -> str: + """Return the configured ACMS vector backend name.""" + + raw = self._config_service.resolve("index.vector.backend").value + return str(raw or "none").strip().lower() + + def acms_enabled(self) -> bool: + """Whether ACMS vector indexing is configured to use FAISS.""" + + return self.acms_backend_name() == "faiss" + + def acms_count(self, *, project: str | None = None) -> int: + """Return the number of indexed ACMS embeddings.""" + + store = self._acms_store or self._load_acms_store() + if store is None: + return 0 + + if project is None: + return len(getattr(store, "index_to_docstore_id", {})) + + count = 0 + for _store_id, metadata in self._iter_store_metadata(store): + if metadata.get("project") == project: + count += 1 + return count + + def acms_index_embedding( + self, + project: str, + doc_id: str, + embedding: list[float], + metadata: dict[str, str], + *, + content: str | None = None, + ) -> None: + """Index one ACMS embedding into the shared FAISS store.""" + + if not project or not project.strip(): + raise ValueError("project must be a non-empty string") + if not doc_id or not doc_id.strip(): + raise ValueError("doc_id must be a non-empty string") + if not embedding: + raise ValueError("embedding must be a non-empty list") + if not self.acms_enabled(): + raise ConfigurationError( + "ACMS vector indexing is disabled. Set index.vector.backend to faiss " + "to enable the FAISS ACMS backend." + ) + + faiss_cls = self._require_faiss() + store = self._acms_store or self._load_acms_store() + stored_metadata = dict(metadata) + stored_metadata["project"] = project + stored_metadata["doc_id"] = doc_id + payload = (content or "").strip() or stored_metadata.get( + "location", f"embedding:{doc_id}" + ) + store_id = self._acms_store_id(project, doc_id) + + if store is None: + store = cast( + _FaissAcmsStoreProtocol, + faiss_cls.from_embeddings( + [(payload, list(embedding))], + self._create_acms_embeddings(), + metadatas=[stored_metadata], + ids=[store_id], + ), + ) + else: + self._delete_store_id(store, store_id) + store.add_embeddings( + [(payload, list(embedding))], + metadatas=[stored_metadata], + ids=[store_id], + ) + + self._acms_store = store + self._save_acms_store() + + def acms_search_by_vector( + self, + embedding: list[float], + *, + top_k: int = 20, + scope: frozenset[str] | None = None, + project: str | None = None, + ) -> list[dict[str, Any]]: + """Search ACMS embeddings by vector similarity.""" + + if not embedding: + raise ValueError("embedding must be a non-empty list") + if top_k < 1: + raise ValueError(f"top_k must be positive, got {top_k}") + if not self.acms_enabled(): + return [] + + store = self._acms_store or self._load_acms_store() + if store is None: + return [] + + scope_set = set(scope or ()) + + def _matches(metadata: dict[str, Any]) -> bool: + if project is not None and metadata.get("project") != project: + return False + return not scope_set or metadata.get("resource_id") in scope_set + + raw_results = store.similarity_search_with_score_by_vector( + list(embedding), + k=top_k, + filter=_matches if (scope_set or project is not None) else None, + fetch_k=max(20, top_k * 5), + ) + results: list[dict[str, Any]] = [] + for document, distance in raw_results: + metadata = { + str(key): str(value) + for key, value in document.metadata.items() + if value is not None + } + results.append( + { + "content": document.page_content, + "distance": float(distance), + "score": self._distance_to_relevance(distance), + "metadata": metadata, + } + ) + return results[:top_k] + + def acms_remove_embedding(self, project: str, doc_id: str) -> None: + """Remove one ACMS embedding from the shared FAISS store.""" + + if not project or not project.strip(): + raise ValueError("project must be a non-empty string") + if not doc_id or not doc_id.strip(): + raise ValueError("doc_id must be a non-empty string") + + store = self._acms_store or self._load_acms_store() + if store is None: + return + + self._delete_store_id(store, self._acms_store_id(project, doc_id)) + if getattr(store, "index_to_docstore_id", {}): + self._acms_store = store + self._save_acms_store() + return + + self._acms_store = None + self._remove_acms_store_files() + # ------------------------------------------------------------------ # Internal helpers # ------------------------------------------------------------------ @@ -113,6 +343,13 @@ class VectorStoreService: "CLEVERAGENTS_VECTOR_STORE_ENABLED to true to enable semantic search." ) + def _require_faiss(self) -> Any: + if FAISS is None: + raise ConfigurationError( + "FAISS library not installed -- vector search will not work" + ) + return FAISS + def _prepare_documents( self, plan_id: int ) -> tuple[list[str], list[dict[str, Any]]]: @@ -148,21 +385,53 @@ class VectorStoreService: f"'{self.settings.vector_embeddings_provider}'" ) + def _create_acms_embeddings(self) -> Embeddings: + provider = ( + str( + self._config_service.resolve("index.embedding.provider").value or "fake" + ) + .strip() + .lower() + ) + model_value = self._config_service.resolve("index.embedding.model").value + dims_value = self._config_service.resolve("index.embedding.dimensions").value + + if provider in {"fake", "local"}: + if dims_value is None: + dims_value = self.settings.vector_embeddings_dimension + return FakeEmbeddings(size=int(dims_value)) + + if provider in {"openai", "azure", "azureopenai"}: + try: + from langchain_openai import OpenAIEmbeddings + except ImportError as exc: # pragma: no cover - import guard + raise ConfigurationError( + "langchain-openai is required for OpenAI embeddings" + ) from exc + + kwargs: dict[str, Any] = {"model": model_value or "text-embedding-3-small"} + if dims_value is not None: + kwargs["dimensions"] = int(dims_value) + return OpenAIEmbeddings(**kwargs) + + raise ConfigurationError(f"Unsupported ACMS embedding provider '{provider}'") + def _plan_store_dir(self, plan_id: int) -> Path: base = self.settings.resolve_vector_store_path() plan_dir = base / f"plan_{plan_id}" plan_dir.mkdir(parents=True, exist_ok=True) return plan_dir - def _load_local_index(self, plan_id: int) -> FAISS | None: + def _load_local_index(self, plan_id: int) -> _FaissPlanStoreProtocol | None: plan_dir = self._plan_store_dir(plan_id) index_path = plan_dir / "index.faiss" store_path = plan_dir / "index.pkl" if not index_path.exists() or not store_path.exists(): return None embeddings = self._create_embeddings() + faiss_cls = self._require_faiss() try: - store = FAISS.load_local( + store = faiss_cls.load_local( str(plan_dir), embeddings, allow_dangerous_deserialization=True, @@ -179,3 +448,61 @@ class VectorStoreService: for path in (index_path, store_path): if path.exists(): path.unlink(missing_ok=True) + + def _acms_store_dir(self) -> Path: + raw_dir = self._config_service.resolve("index.vector.dir").value + base = Path(str(raw_dir)).expanduser() + if not base.is_absolute(): + base = Path.cwd() / base + base.mkdir(parents=True, exist_ok=True) + return base.resolve() + + def _load_acms_store(self) -> _FaissAcmsStoreProtocol | None: + store_dir = self._acms_store_dir() + index_path = store_dir / "index.faiss" + store_path = store_dir / "index.pkl" + if not index_path.exists() or not store_path.exists(): + return None + + faiss_cls = self._require_faiss() + try: + store = faiss_cls.load_local( + str(store_dir), + self._create_acms_embeddings(), + allow_dangerous_deserialization=True, + ) + except (FileNotFoundError, ValueError): + return None + self._acms_store = store + return store + + def _save_acms_store(self) -> None: + if self._acms_store is None: + return + self._acms_store.save_local(str(self._acms_store_dir())) + + def _remove_acms_store_files(self) -> None: + store_dir = self._acms_store_dir() + for path in (store_dir / "index.faiss", store_dir / "index.pkl"): + if path.exists(): + path.unlink(missing_ok=True) + + def _acms_store_id(self, project: str, doc_id: str) -> str: + return f"{project}::{doc_id}" + + def _delete_store_id(self, store: _FaissAcmsStoreProtocol, store_id: str) -> None: + docstore_ids = set(getattr(store, "index_to_docstore_id", {}).values()) + if store_id in docstore_ids: + store.delete(ids=[store_id]) + + def _iter_store_metadata( + self, store: _FaissAcmsStoreProtocol + ) -> Iterable[tuple[str, dict[str, str]]]: + for store_id in getattr(store, "index_to_docstore_id", {}).values(): + document = store.docstore.search(store_id) + metadata = getattr(document, "metadata", {}) + yield store_id, {str(key): str(value) for key, value in metadata.items()} + + def _distance_to_relevance(self, distance: Any) -> float: + numeric_distance = max(float(distance), 0.0) + return 1.0 / (1.0 + numeric_distance)