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