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

1"""Vector-backed memory backend — semantic search via DocumentVectorStoreProtocol.""" 

2 

3from __future__ import annotations 

4 

5from datetime import UTC, datetime 

6from typing import TYPE_CHECKING, Any, cast 

7 

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 

21 

22if TYPE_CHECKING: 

23 from collections.abc import Awaitable, Callable 

24 

25 

26class VectorMemoryBackend: 

27 """MemoryStoreProtocol that persists entries as vector-searchable documents.""" 

28 

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 

40 

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) 

45 

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 ) 

62 

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) 

73 

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") 

91 

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) 

116 

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() 

132 

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] 

136 

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) 

141 

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) 

151 

152 async def clear(self, owner_id: str) -> None: 

153 if self._fallback is not None: 

154 await self._fallback.clear(owner_id) 

155 

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 ) 

162 

163 

164__all__ = ["VectorMemoryBackend"]