Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-memory/src/lexigram/ai/memory/backends/vector.py: 30%
63 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-25 07:19 +0800
1"""Vector-backed memory backend — semantic search via DocumentVectorStoreProtocol."""
3from __future__ import annotations
5from datetime import UTC, datetime
6from typing import TYPE_CHECKING, Any, cast
8from lexigram.ai.memory.exceptions import EmbeddingError, MemoryStoreError
9from lexigram.contracts.ai.memory import (
10 MemoryEntry,
11 MemoryQuery,
12 MemorySearchResult,
13 MemoryStoreProtocol,
14)
15from lexigram.contracts.ai.vector import (
16 Document,
17 DocumentVectorStoreProtocol,
18 SearchResultProtocol,
19)
20from lexigram.contracts.core import HealthCheckResult, HealthStatus
22if TYPE_CHECKING:
23 from collections.abc import Awaitable, Callable
26class VectorMemoryBackend:
27 """MemoryStoreProtocol that persists entries as vector-searchable documents."""
29 def __init__(
30 self,
31 vector_store: DocumentVectorStoreProtocol,
32 embed_fn: Callable[[str], Awaitable[list[float]]] | None = None,
33 collection: str = "memory",
34 fallback: MemoryStoreProtocol | None = None,
35 ) -> None:
36 self._vs = vector_store
37 self._embed_fn = embed_fn
38 self._collection = collection
39 self._fallback = fallback
41 async def _embed(self, text: str) -> list[float]:
42 if self._embed_fn is None:
43 raise EmbeddingError("No embed_fn provided — cannot produce embedding")
44 return await self._embed_fn(text)
46 def _to_document(self, entry: MemoryEntry) -> Document:
47 metadata: dict[str, Any] = {
48 "owner_id": entry.owner_id,
49 "content": entry.content,
50 "role": entry.role,
51 "timestamp": entry.timestamp.isoformat(),
52 "importance": entry.importance,
53 "memory_metadata": entry.metadata,
54 "collection": self._collection,
55 }
56 return Document(
57 id=entry.id,
58 text=entry.content,
59 metadata=metadata,
60 embedding=entry.embedding,
61 )
63 def _to_memory_result(
64 self,
65 hit: SearchResultProtocol,
66 query: MemoryQuery,
67 ) -> MemorySearchResult:
68 metadata = hit.document.metadata
69 timestamp_raw = metadata.get("timestamp", datetime.now(UTC).isoformat())
70 timestamp = datetime.fromisoformat(str(timestamp_raw))
71 if timestamp.tzinfo is None:
72 timestamp = timestamp.replace(tzinfo=UTC)
74 entry = MemoryEntry(
75 id=hit.document.id or "",
76 owner_id=str(metadata.get("owner_id", "")),
77 content=metadata.get("content", hit.document.text),
78 role=str(metadata.get("role", "user")),
79 timestamp=timestamp,
80 importance=float(metadata.get("importance", 0.5)),
81 metadata=cast("dict[str, Any]", metadata.get("memory_metadata", {})),
82 embedding=cast(
83 "list[float] | None", getattr(hit.document, "embedding", None)
84 ),
85 )
86 combined_score = (
87 query.relevance_weight * float(hit.score)
88 + query.importance_weight * entry.importance
89 )
90 return MemorySearchResult(entry=entry, score=combined_score, source="vector")
92 async def store(self, entry: MemoryEntry) -> None:
93 vector = entry.embedding
94 if vector is None:
95 vector = await self._embed(entry.content)
96 doc = self._to_document(
97 MemoryEntry(
98 id=entry.id,
99 owner_id=entry.owner_id,
100 content=entry.content,
101 role=entry.role,
102 timestamp=entry.timestamp,
103 importance=entry.importance,
104 metadata=entry.metadata,
105 embedding=vector,
106 )
107 )
108 add_result = await self._vs.add([doc])
109 if add_result.is_err():
110 raise MemoryStoreError(
111 f"Vector add failed: {entry.id}",
112 store="vector",
113 ) from add_result.unwrap_err()
114 if self._fallback is not None:
115 await self._fallback.store(entry)
117 async def retrieve(self, query: MemoryQuery) -> list[MemorySearchResult]:
118 query_vector = await self._embed(query.query)
119 search_result = await self._vs.search(
120 query_vector,
121 top_k=query.top_k,
122 filters={
123 "collection": self._collection,
124 "owner_id": query.owner_id,
125 **(query.filters or {}),
126 },
127 )
128 if search_result.is_err():
129 raise MemoryStoreError(
130 "Vector search failed", store="vector"
131 ) from search_result.unwrap_err()
133 hits = search_result.unwrap_or([])
134 results = [self._to_memory_result(hit, query) for hit in hits]
135 return [result for result in results if result.score >= query.min_relevance]
137 async def get_recent(self, n: int, owner_id: str) -> list[MemoryEntry]:
138 if self._fallback is None:
139 return []
140 return await self._fallback.get_recent(n, owner_id)
142 async def delete(self, entry_id: str, owner_id: str) -> None:
143 delete_result = await self._vs.delete([entry_id])
144 if delete_result.is_err():
145 raise MemoryStoreError(
146 f"Vector delete failed: {entry_id}",
147 store="vector",
148 ) from delete_result.unwrap_err()
149 if self._fallback is not None:
150 await self._fallback.delete(entry_id, owner_id)
152 async def clear(self, owner_id: str) -> None:
153 if self._fallback is not None:
154 await self._fallback.clear(owner_id)
156 async def health_check(self, timeout: float = 5.0) -> HealthCheckResult:
157 return HealthCheckResult(
158 component="memory.vector",
159 status=HealthStatus.HEALTHY,
160 details={"collection": self._collection, "timeout": timeout},
161 )
164__all__ = ["VectorMemoryBackend"]