From eff446f5e8165fa2a0f2862e7411346b079f53ce Mon Sep 17 00:00:00 2001 From: Aditya Chhabra Date: Thu, 26 Mar 2026 13:13:35 +0000 Subject: [PATCH] feat(acms): integrate FAISS into ACMS vector backend protocol Add FAISS-backed ACMS read and write adapters on top of the shared VectorStoreService, wire them through the DI container, and cover indexing, scoped search, removal, and benchmark behavior with Behave and ASV. Address review feedback by keeping the commit scoped to the FAISS backend work and replacing the loose vector-store cache typing with explicit FAISS store protocols instead of Any. ISSUES CLOSED: #871 --- benchmarks/faiss_vector_backend_bench.py | 102 ++++++ features/faiss_vector_backend.feature | 46 +++ features/steps/faiss_vector_backend_steps.py | 338 +++++++++++++++++ src/cleveragents/application/container.py | 29 +- .../services/faiss_vector_backend.py | 204 +++++++++++ .../services/vector_store_service.py | 343 +++++++++++++++++- 6 files changed, 1048 insertions(+), 14 deletions(-) create mode 100644 benchmarks/faiss_vector_backend_bench.py create mode 100644 features/faiss_vector_backend.feature create mode 100644 features/steps/faiss_vector_backend_steps.py create mode 100644 src/cleveragents/application/services/faiss_vector_backend.py 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) -- 2.52.0