feat(acms): integrate FAISS into ACMS vector backend protocol #1165

Merged
aditya merged 1 commits from feature/m6-faiss-acms-backend into master 2026-03-30 13:40:25 +00:00
6 changed files with 1048 additions and 14 deletions
+102
View File
@@ -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,
)
+46
View File
@@ -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
@@ -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
)
+23 -6
View File
@@ -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,
@@ -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",
]
@@ -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)