Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-feedback/src/lexigram/ai/feedback/storage/database.py: 36%

73 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1"""SQL-backed feedback store using :class:`~lexigram.contracts.data.DatabaseProviderProtocol`. 

2 

3Records all four feedback types, extracts ``session_id`` from item context 

4into a dedicated column for efficient session-scoped queries, and serialises 

5``value``/``context``/``metadata`` as JSON strings. 

6""" 

7 

8from __future__ import annotations 

9 

10from datetime import UTC, datetime, timedelta 

11from typing import TYPE_CHECKING, Any 

12 

13from lexigram.ai.feedback.exceptions import FeedbackError 

14from lexigram.ai.feedback.storage.protocols import FeedbackSummary 

15from lexigram.ai.feedback.types import FeedbackItem, FeedbackType 

16from lexigram.logging import ( 

17 get_logger, 

18) 

19from lexigram.result import Err, Ok, Result 

20from lexigram.serialization import dumps_str, loads_str 

21 

22if TYPE_CHECKING: 

23 from lexigram.contracts.data import DatabaseProviderProtocol 

24 

25logger = get_logger(__name__) 

26 

27_CREATE_TABLE = """ 

28CREATE TABLE IF NOT EXISTS ai_feedback ( 

29 id TEXT NOT NULL PRIMARY KEY, 

30 type TEXT NOT NULL, 

31 value TEXT NOT NULL, 

32 context TEXT NOT NULL DEFAULT '{}', 

33 metadata TEXT NOT NULL DEFAULT '{}', 

34 session_id TEXT, 

35 owner_id TEXT NOT NULL, 

36 created_at TEXT NOT NULL 

37) 

38""" 

39 

40_CREATE_IDX_OWNER = ( 

41 "CREATE INDEX IF NOT EXISTS idx_ai_feedback_owner ON ai_feedback (owner_id)" 

42) 

43_CREATE_IDX_SESSION = ( 

44 "CREATE INDEX IF NOT EXISTS idx_ai_feedback_session ON ai_feedback (session_id)" 

45) 

46_CREATE_IDX_TYPE = ( 

47 "CREATE INDEX IF NOT EXISTS idx_ai_feedback_type ON ai_feedback (type)" 

48) 

49_CREATE_IDX_CREATED = ( 

50 "CREATE INDEX IF NOT EXISTS idx_ai_feedback_created_at ON ai_feedback (created_at)" 

51) 

52 

53_INSERT = ( 

54 "INSERT OR IGNORE INTO ai_feedback " 

55 "(id, type, value, context, metadata, session_id, owner_id, created_at) " 

56 "VALUES (?, ?, ?, ?, ?, ?, ?, ?)" 

57) 

58 

59 

60class DatabaseFeedbackStore: 

61 """Persist and query feedback items via :class:`~lexigram.contracts.data.DatabaseProviderProtocol`. 

62 

63 The backing table (``ai_feedback``) is created lazily on first write. 

64 Indexed columns: ``session_id``, ``type``, ``created_at``. 

65 

66 Args: 

67 provider: DI-injected database provider. 

68 """ 

69 

70 def __init__(self, provider: DatabaseProviderProtocol) -> None: 

71 self._db = provider 

72 self._initialised = False 

73 

74 # ------------------------------------------------------------------ 

75 # Lifecycle 

76 # ------------------------------------------------------------------ 

77 

78 async def _ensure_table(self) -> None: 

79 """Create backing table + indexes on first use.""" 

80 if self._initialised: 

81 return 

82 await self._db.execute(_CREATE_TABLE) 

83 await self._db.execute(_CREATE_IDX_OWNER) 

84 await self._db.execute(_CREATE_IDX_SESSION) 

85 await self._db.execute(_CREATE_IDX_TYPE) 

86 await self._db.execute(_CREATE_IDX_CREATED) 

87 self._initialised = True 

88 

89 # ------------------------------------------------------------------ 

90 # FeedbackStoreProtocol implementation 

91 # ------------------------------------------------------------------ 

92 

93 async def save(self, feedback: FeedbackItem) -> Result[str, FeedbackError]: 

94 """Persist *feedback* and return its ID on success. 

95 

96 Args: 

97 feedback: Item to store. 

98 

99 Returns: 

100 ``Ok(feedback.id)`` on success, ``Err(FeedbackError)`` on failure. 

101 """ 

102 try: 

103 await self._ensure_table() 

104 session_id: str | None = feedback.context.get("session_id") 

105 await self._db.execute( 

106 _INSERT, 

107 [ 

108 feedback.id, 

109 feedback.feedback_type.value, 

110 dumps_str(feedback.value), 

111 dumps_str(feedback.context), 

112 dumps_str(feedback.metadata), 

113 session_id, 

114 feedback.owner_id, 

115 feedback.created_at.isoformat(), 

116 ], 

117 ) 

118 return Ok(feedback.id) 

119 except (ConnectionError, TimeoutError, OSError) as exc: 

120 logger.error( 

121 "feedback_save_failed", feedback_id=feedback.id, error=str(exc) 

122 ) 

123 return Err(FeedbackError(f"Failed to save feedback {feedback.id}: {exc}")) 

124 

125 async def find_by_session( 

126 self, session_id: str, *, owner_id: str 

127 ) -> list[FeedbackItem]: 

128 """Return all items collected during *session_id* for *owner_id*, newest first. 

129 

130 Args: 

131 session_id: Session identifier to filter by. 

132 owner_id: Owner scope; only this owner's items are returned. 

133 

134 Returns: 

135 Matching feedback items ordered by creation time descending. 

136 """ 

137 await self._ensure_table() 

138 result = await self._db.execute_query( 

139 "SELECT * FROM ai_feedback WHERE session_id = ? AND owner_id = ? " 

140 "ORDER BY created_at DESC", 

141 [session_id, owner_id], 

142 ) 

143 return [self._row_to_item(row) for row in result.rows] 

144 

145 async def find_by_type( 

146 self, 

147 feedback_type: FeedbackType, 

148 *, 

149 owner_id: str, 

150 limit: int = 100, 

151 ) -> list[FeedbackItem]: 

152 """Return feedback items of *feedback_type* for *owner_id*, newest first. 

153 

154 Args: 

155 feedback_type: Type to filter by. 

156 owner_id: Owner scope; only this owner's items are returned. 

157 limit: Maximum number of results (default: 100). 

158 

159 Returns: 

160 Matching feedback items ordered by creation time descending. 

161 """ 

162 await self._ensure_table() 

163 result = await self._db.execute_query( 

164 "SELECT * FROM ai_feedback WHERE type = ? AND owner_id = ? " 

165 "ORDER BY created_at DESC LIMIT ?", 

166 [feedback_type.value, owner_id, limit], 

167 ) 

168 return [self._row_to_item(row) for row in result.rows] 

169 

170 async def aggregate( 

171 self, *, owner_id: str, window_hours: int = 24 

172 ) -> FeedbackSummary: 

173 """Compute summary statistics for *owner_id* in the last *window_hours* hours. 

174 

175 Args: 

176 owner_id: Owner scope; only this owner's items are aggregated. 

177 window_hours: Look-back window in hours (default: 24). 

178 

179 Returns: 

180 Aggregated :class:`FeedbackSummary` for the window. 

181 """ 

182 await self._ensure_table() 

183 since = (datetime.now(UTC) - timedelta(hours=window_hours)).isoformat() 

184 

185 # Total count 

186 count_result = await self._db.execute_query( 

187 "SELECT COUNT(*) AS cnt FROM ai_feedback " 

188 "WHERE created_at >= ? AND owner_id = ?", 

189 [since, owner_id], 

190 ) 

191 total = int((count_result.rows[0] or {}).get("cnt", 0)) 

192 

193 # Average rating 

194 avg_result = await self._db.execute_query( 

195 "SELECT AVG(CAST(value AS REAL)) AS avg_rating " 

196 "FROM ai_feedback WHERE type = ? AND created_at >= ? AND owner_id = ?", 

197 [FeedbackType.RATING.value, since, owner_id], 

198 ) 

199 raw_avg = (avg_result.rows[0] or {}).get("avg_rating") 

200 average_rating = float(raw_avg) if raw_avg is not None else None 

201 

202 # Count by type 

203 type_result = await self._db.execute_query( 

204 "SELECT type, COUNT(*) AS cnt FROM ai_feedback " 

205 "WHERE created_at >= ? AND owner_id = ? GROUP BY type", 

206 [since, owner_id], 

207 ) 

208 count_by_type = {row["type"]: int(row["cnt"]) for row in type_result.rows} 

209 

210 return FeedbackSummary( 

211 total_count=total, 

212 average_rating=average_rating, 

213 count_by_type=count_by_type, 

214 ) 

215 

216 # ------------------------------------------------------------------ 

217 # Private helpers 

218 # ------------------------------------------------------------------ 

219 

220 @staticmethod 

221 def _row_to_item(row: Any) -> FeedbackItem: 

222 """Deserialise a database row back into a :class:`FeedbackItem`. 

223 

224 Args: 

225 row: Raw row dict from the database. 

226 

227 Returns: 

228 Reconstructed :class:`FeedbackItem`. 

229 """ 

230 try: 

231 value: Any = loads_str(row["value"]) 

232 except (ValueError, TypeError): 

233 value = row["value"] 

234 

235 try: 

236 context: dict[str, Any] = loads_str(row["context"]) or {} 

237 except (ValueError, TypeError): 

238 context = {} 

239 

240 try: 

241 metadata: dict[str, Any] = loads_str(row["metadata"]) or {} 

242 except (ValueError, TypeError): 

243 metadata = {} 

244 

245 return FeedbackItem( 

246 feedback_type=FeedbackType(row["type"]), 

247 value=value, 

248 owner_id=str(row["owner_id"]), 

249 context=context, 

250 metadata=metadata, 

251 id=row["id"], 

252 created_at=datetime.fromisoformat(row["created_at"]), 

253 ) 

254 

255 

256__all__ = ["DatabaseFeedbackStore"]