Coverage for src/lexigram/notification/inbox/database.py: 32%

80 statements  

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

1"""Database-persisted inbox store implementation. 

2 

3Uses :class:`~lexigram.contracts.data.sql.database.DatabaseProviderProtocol` 

4so the underlying engine (PostgreSQL, MySQL, SQLite …) is fully decoupled. 

5""" 

6 

7from __future__ import annotations 

8 

9from datetime import UTC, datetime 

10from typing import Any 

11 

12from lexigram.contracts.core import HealthCheckResult, HealthStatus 

13from lexigram.contracts.data.sql.database import DatabaseProviderProtocol 

14from lexigram.contracts.exceptions import DatabaseError 

15from lexigram.contracts.notification.inbox import InboxMessage 

16from lexigram.logging import get_logger 

17 

18logger = get_logger(__name__) 

19 

20_DEFAULT_TABLE = "notification_inbox_messages" 

21 

22_CREATE_TABLE_SQL = """ 

23CREATE TABLE IF NOT EXISTS {table} ( 

24 id TEXT PRIMARY KEY, 

25 user_id TEXT NOT NULL, 

26 title TEXT NOT NULL, 

27 body TEXT NOT NULL, 

28 read BOOLEAN NOT NULL DEFAULT FALSE, 

29 created_at TIMESTAMPTZ NOT NULL, 

30 metadata JSONB NOT NULL DEFAULT '{{}}' 

31); 

32CREATE INDEX IF NOT EXISTS {table}_user_created_idx 

33 ON {table} (user_id, created_at DESC); 

34""" 

35 

36 

37class DatabaseInboxStore: 

38 """SQL-backed :class:`InboxStoreProtocol` implementation. 

39 

40 Reads and writes to a single flat table. The table is created 

41 automatically on first use if it does not exist:: 

42 

43 CREATE TABLE notification_inbox_messages ( 

44 id TEXT PRIMARY KEY, 

45 user_id TEXT NOT NULL, 

46 title TEXT NOT NULL, 

47 body TEXT NOT NULL, 

48 read BOOLEAN NOT NULL DEFAULT FALSE, 

49 created_at TIMESTAMPTZ NOT NULL, 

50 metadata JSONB NOT NULL DEFAULT '{}' 

51 ); 

52 CREATE INDEX ON notification_inbox_messages (user_id, created_at DESC); 

53 

54 Args: 

55 db: Database provider injected via DI. 

56 table: Table name. Defaults to ``notification_inbox_messages``. 

57 """ 

58 

59 def __init__( 

60 self, 

61 db: DatabaseProviderProtocol, 

62 *, 

63 table: str = _DEFAULT_TABLE, 

64 ) -> None: 

65 self._db = db 

66 self._table = table 

67 

68 # ------------------------------------------------------------------ 

69 # InboxStoreProtocol implementation 

70 # ------------------------------------------------------------------ 

71 

72 async def _ensure_table(self) -> None: 

73 """Create the inbox table (and index) if it does not exist.""" 

74 sql = _CREATE_TABLE_SQL.format(table=self._table) 

75 for statement in sql.split(";"): 

76 stmt = statement.strip() 

77 if stmt: 

78 await self._db.execute_query(stmt) 

79 

80 async def save(self, message: InboxMessage) -> None: 

81 """Persist *message* in the database. 

82 

83 Args: 

84 message: The inbox message to store. 

85 """ 

86 await self._ensure_table() 

87 await self._db.execute_insert( 

88 self._table, 

89 { 

90 "id": message.id, 

91 "user_id": message.user_id, 

92 "title": message.title, 

93 "body": message.body, 

94 "read": message.read, 

95 "created_at": message.created_at, 

96 "metadata": message.metadata, 

97 }, 

98 ) 

99 logger.debug("inbox.db.saved", message_id=message.id, user_id=message.user_id) 

100 

101 async def get(self, message_id: str) -> InboxMessage | None: 

102 """Return the message with *message_id*, or ``None``. 

103 

104 Args: 

105 message_id: ID to look up. 

106 """ 

107 await self._ensure_table() 

108 result = await self._db.execute_query( 

109 f"SELECT * FROM {self._table} WHERE id = $1", # noqa: S608 

110 [message_id], 

111 ) 

112 if not result.rows: 

113 return None 

114 return self._row_to_message(result.rows[0]) 

115 

116 async def list_for_user( 

117 self, 

118 user_id: str, 

119 *, 

120 unread_only: bool = False, 

121 ) -> list[InboxMessage]: 

122 """Return messages for *user_id* in reverse-chronological order. 

123 

124 Args: 

125 user_id: Filter to this user. 

126 unread_only: When ``True`` only unread messages are returned. 

127 """ 

128 await self._ensure_table() 

129 if unread_only: 

130 result = await self._db.execute_query( 

131 f"SELECT * FROM {self._table}" # noqa: S608 

132 " WHERE user_id = $1 AND read = FALSE" 

133 " ORDER BY created_at DESC", 

134 [user_id], 

135 ) 

136 else: 

137 result = await self._db.execute_query( 

138 f"SELECT * FROM {self._table}" # noqa: S608 

139 " WHERE user_id = $1" 

140 " ORDER BY created_at DESC", 

141 [user_id], 

142 ) 

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

144 

145 async def mark_read(self, message_id: str, user_id: str) -> None: 

146 """Mark *message_id* as read, guarded by *user_id* ownership. 

147 

148 Args: 

149 message_id: ID of the message to mark. 

150 user_id: Expected owner. 

151 """ 

152 await self._ensure_table() 

153 await self._db.execute_update( 

154 self._table, 

155 {"read": True}, 

156 "id = $1 AND user_id = $2", 

157 [message_id, user_id], 

158 ) 

159 logger.debug("inbox.db.marked_read", message_id=message_id, user_id=user_id) 

160 

161 async def mark_all_read(self, user_id: str) -> None: 

162 """Mark every unread message for *user_id* as read. 

163 

164 Args: 

165 user_id: Target user. 

166 """ 

167 await self._ensure_table() 

168 await self._db.execute_update( 

169 self._table, 

170 {"read": True}, 

171 "user_id = $1 AND read = FALSE", 

172 [user_id], 

173 ) 

174 logger.debug("inbox.db.all_marked_read", user_id=user_id) 

175 

176 async def delete(self, message_id: str, user_id: str) -> None: 

177 """Delete *message_id*, guarded by *user_id* ownership. 

178 

179 Args: 

180 message_id: ID to remove. 

181 user_id: Expected owner. 

182 """ 

183 await self._ensure_table() 

184 await self._db.execute_delete( 

185 self._table, 

186 "id = $1 AND user_id = $2", 

187 [message_id, user_id], 

188 ) 

189 logger.debug("inbox.db.deleted", message_id=message_id, user_id=user_id) 

190 

191 async def count_unread(self, user_id: str) -> int: 

192 """Return unread message count for *user_id*. 

193 

194 Args: 

195 user_id: Target user. 

196 """ 

197 await self._ensure_table() 

198 result = await self._db.execute_query( 

199 f"SELECT COUNT(*) AS cnt FROM {self._table}" # noqa: S608 

200 " WHERE user_id = $1 AND read = FALSE", 

201 [user_id], 

202 ) 

203 if not result.rows: 

204 return 0 

205 return int(result.rows[0].get("cnt", 0)) 

206 

207 async def clear_all(self, user_id: str) -> int: 

208 """Delete all messages for *user_id*. 

209 

210 Args: 

211 user_id: Target user. 

212 

213 Returns: 

214 Number of messages deleted. 

215 """ 

216 await self._ensure_table() 

217 count_result = await self._db.execute_query( 

218 f"SELECT COUNT(*) AS cnt FROM {self._table}" # noqa: S608 

219 " WHERE user_id = $1", 

220 [user_id], 

221 ) 

222 count = int(count_result.rows[0].get("cnt", 0)) if count_result.rows else 0 

223 await self._db.execute_delete( 

224 self._table, 

225 "user_id = $1", 

226 [user_id], 

227 ) 

228 logger.debug("inbox.db.cleared_all", user_id=user_id, count=count) 

229 return count 

230 

231 async def health_check(self, timeout: float = 5.0) -> HealthCheckResult: 

232 """Run a lightweight read-only database probe for inbox storage.""" 

233 try: 

234 await self._db.execute_query( 

235 f"SELECT 1 FROM {self._table} LIMIT 1", # noqa: S608 

236 [], 

237 ) 

238 except ( 

239 DatabaseError, 

240 ConnectionError, 

241 RuntimeError, 

242 TimeoutError, 

243 ValueError, 

244 ) as exc: 

245 return HealthCheckResult( 

246 component="inbox_store", 

247 status=HealthStatus.UNHEALTHY, 

248 error=str(exc), 

249 details={"backend": "database", "table": self._table}, 

250 ) 

251 

252 return HealthCheckResult( 

253 component="inbox_store", 

254 status=HealthStatus.HEALTHY, 

255 details={"backend": "database", "table": self._table}, 

256 ) 

257 

258 # ------------------------------------------------------------------ 

259 # Internal helpers 

260 # ------------------------------------------------------------------ 

261 

262 @staticmethod 

263 def _row_to_message(row: dict[str, Any]) -> InboxMessage: 

264 """Convert a raw DB row dict to an :class:`InboxMessage`. 

265 

266 Args: 

267 row: Raw row returned by the database driver. 

268 

269 Returns: 

270 Hydrated :class:`InboxMessage` instance. 

271 """ 

272 created_at = row["created_at"] 

273 if isinstance(created_at, str): 

274 created_at = datetime.fromisoformat(created_at) 

275 elif not isinstance(created_at, datetime): 

276 created_at = datetime.now(UTC) 

277 

278 metadata = row.get("metadata") or {} 

279 if not isinstance(metadata, dict): 

280 metadata = {} 

281 

282 return InboxMessage( 

283 id=str(row["id"]), 

284 user_id=str(row["user_id"]), 

285 title=str(row["title"]), 

286 body=str(row["body"]), 

287 read=bool(row.get("read", False)), 

288 created_at=created_at, 

289 metadata=metadata, 

290 ) 

291 

292 

293__all__ = ["DatabaseInboxStore"]