Coverage for src/lexigram/notification/inbox/memory.py: 40%

47 statements  

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

1"""In-memory inbox store implementation. 

2 

3Suitable for unit tests and single-process deployments that do not 

4require durable notification persistence. 

5""" 

6 

7from __future__ import annotations 

8 

9from dataclasses import replace 

10from datetime import UTC, datetime 

11 

12from lexigram.contracts.core import HealthCheckResult, HealthStatus 

13from lexigram.contracts.notification.inbox import InboxMessage 

14from lexigram.logging import get_logger 

15 

16logger = get_logger(__name__) 

17 

18 

19class InMemoryInboxStore: 

20 """In-process :class:`InboxStoreProtocol` implementation backed by a dict. 

21 

22 Thread-safety note: this store is *not* thread-safe. Within a single 

23 async event loop it is safe; for multi-threaded scenarios inject a 

24 SQL-backed :class:`DatabaseInboxStore` instead. 

25 """ 

26 

27 def __init__(self) -> None: 

28 self._store: dict[str, InboxMessage] = {} 

29 

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

31 """Persist *message* in memory. 

32 

33 Args: 

34 message: The inbox message to store. 

35 """ 

36 self._store[message.id] = message 

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

38 

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

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

41 

42 Args: 

43 message_id: ID to look up. 

44 """ 

45 return self._store.get(message_id) 

46 

47 async def list_for_user( 

48 self, 

49 user_id: str, 

50 *, 

51 unread_only: bool = False, 

52 ) -> list[InboxMessage]: 

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

54 

55 Args: 

56 user_id: Filter to this user. 

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

58 """ 

59 results = [ 

60 m 

61 for m in self._store.values() 

62 if m.user_id == user_id and (not unread_only or not m.read) 

63 ] 

64 results.sort(key=lambda m: m.created_at, reverse=True) 

65 return results 

66 

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

68 """Mark *message_id* as read if it belongs to *user_id*. 

69 

70 ``InboxMessage`` is frozen; we create a replacement with 

71 ``read=True`` and update the store in-place. 

72 

73 Args: 

74 message_id: ID of the message to mark. 

75 user_id: Expected owner. 

76 """ 

77 msg = self._store.get(message_id) 

78 if msg is not None and msg.user_id == user_id and not msg.read: 

79 self._store[message_id] = replace(msg, read=True) 

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

81 

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

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

84 

85 Args: 

86 user_id: Target user. 

87 """ 

88 now = datetime.now(UTC) 

89 for mid, msg in list(self._store.items()): 

90 if msg.user_id == user_id and not msg.read: 

91 self._store[mid] = replace(msg, read=True) 

92 logger.debug("inbox.all_marked_read", user_id=user_id, ts=now.isoformat()) 

93 

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

95 """Delete *message_id* if it belongs to *user_id*. 

96 

97 Args: 

98 message_id: ID to remove. 

99 user_id: Expected owner. 

100 """ 

101 msg = self._store.get(message_id) 

102 if msg is not None and msg.user_id == user_id: 

103 del self._store[message_id] 

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

105 

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

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

108 

109 Args: 

110 user_id: Target user. 

111 """ 

112 return sum( 

113 1 for m in self._store.values() if m.user_id == user_id and not m.read 

114 ) 

115 

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

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

118 

119 Args: 

120 user_id: Target user. 

121 

122 Returns: 

123 Number of messages deleted. 

124 """ 

125 ids = [mid for mid, m in self._store.items() if m.user_id == user_id] 

126 for mid in ids: 

127 del self._store[mid] 

128 logger.debug("inbox.cleared_all", user_id=user_id, count=len(ids)) 

129 return len(ids) 

130 

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

132 """Return health for the in-memory inbox backend.""" 

133 _ = timeout 

134 return HealthCheckResult( 

135 component="inbox_store", 

136 status=HealthStatus.HEALTHY, 

137 details={"backend": "memory", "message_count": len(self._store)}, 

138 ) 

139 

140 

141__all__ = ["InMemoryInboxStore"]