Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-governance/src/lexigram/ai/governance/audit/query.py: 30%

64 statements  

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

1"""Governance audit query and streaming export service. 

2 

3Provides rich query, aggregation, and streaming CSV/JSON export of 

4:class:`~lexigram.ai.governance.audit.models.AIAuditEvent` records stored 

5via any :class:`~lexigram.ai.governance.audit.store.AIAuditStore` implementation. 

6""" 

7 

8from __future__ import annotations 

9 

10import csv 

11from datetime import UTC, datetime, timedelta 

12import io 

13from typing import TYPE_CHECKING 

14 

15from lexigram.ai.governance.audit.models import AIAuditEvent, AuditQuery, AuditSummary 

16from lexigram.logging import ( 

17 get_logger, 

18) 

19from lexigram.serialization import dumps_str 

20 

21if TYPE_CHECKING: 

22 from collections.abc import AsyncIterator 

23 

24 from lexigram.ai.governance.audit.store import AIAuditStore 

25 

26logger = get_logger(__name__) 

27 

28_CSV_FIELDS = ( 

29 "event_id", 

30 "event_type", 

31 "timestamp", 

32 "model", 

33 "provider", 

34 "user_id", 

35 "status", 

36 "tokens", 

37 "cost", 

38 "latency_ms", 

39 "metadata", 

40) 

41 

42_EXPORT_PAGE_SIZE = 500 

43 

44 

45def _event_to_csv_row(event: AIAuditEvent) -> list[str]: 

46 """Serialise *event* to a flat list of strings suitable for a CSV row. 

47 

48 Args: 

49 event: The audit event to serialise. 

50 

51 Returns: 

52 Ordered list of string values matching :data:`_CSV_FIELDS`. 

53 """ 

54 return [ 

55 event.event_id, 

56 event.event_type.value, 

57 event.timestamp.isoformat(), 

58 event.model or "", 

59 event.provider or "", 

60 event.user_id or "", 

61 event.status, 

62 str(event.tokens) if event.tokens is not None else "", 

63 str(event.cost) if event.cost is not None else "", 

64 str(event.latency_ms) if event.latency_ms is not None else "", 

65 dumps_str(event.metadata), 

66 ] 

67 

68 

69class AuditQueryService: 

70 """Query and export governance audit records. 

71 

72 All heavy pagination is handled internally so callers can use simple 

73 async-for loops over streaming export methods. 

74 

75 Args: 

76 store: Audit store to delegate persistence operations to. 

77 """ 

78 

79 def __init__(self, store: AIAuditStore) -> None: 

80 self._store = store 

81 

82 async def query( 

83 self, 

84 *, 

85 model: str | None = None, 

86 tenant_id: str | None = None, 

87 start: datetime | None = None, 

88 end: datetime | None = None, 

89 limit: int = 100, 

90 ) -> list[AIAuditEvent]: 

91 """Retrieve audit events matching the given filters. 

92 

93 Args: 

94 model: Restrict to events for this model identifier. 

95 tenant_id: Restrict to events for this user/tenant ID. 

96 start: Inclusive lower bound on event timestamp (UTC). 

97 end: Inclusive upper bound on event timestamp (UTC). 

98 limit: Maximum number of results (default: 100, max: 1 000). 

99 

100 Returns: 

101 Matching events ordered by timestamp descending. 

102 """ 

103 q = AuditQuery( 

104 model=model, 

105 user_id=tenant_id, 

106 start=start, 

107 end=end, 

108 limit=min(limit, 1000), 

109 ) 

110 logger.debug("audit_query", model=model, tenant_id=tenant_id, limit=limit) 

111 return await self._store.query(q) 

112 

113 async def export_csv(self, query: AuditQuery) -> AsyncIterator[bytes]: 

114 """Stream audit records as CSV-encoded bytes. 

115 

116 The first chunk contains the header row. Subsequent chunks are 

117 batches of :data:`_EXPORT_PAGE_SIZE` rows. All pages are fetched 

118 from the underlying store on demand so memory usage stays bounded. 

119 

120 Args: 

121 query: Filter criteria; ``limit`` and ``offset`` are used for 

122 internal pagination and should be left at their defaults. 

123 

124 Yields: 

125 UTF-8 encoded bytes for each batch (header + rows). 

126 """ 

127 # Yield header 

128 buf = io.StringIO() 

129 writer = csv.writer(buf) 

130 writer.writerow(_CSV_FIELDS) 

131 yield buf.getvalue().encode() 

132 

133 offset = 0 

134 page_query = AuditQuery( 

135 start=query.start, 

136 end=query.end, 

137 event_types=query.event_types, 

138 user_id=query.user_id, 

139 model=query.model, 

140 provider=query.provider, 

141 status=query.status, 

142 limit=_EXPORT_PAGE_SIZE, 

143 offset=offset, 

144 ) 

145 

146 while True: 

147 page_query.offset = offset 

148 events = await self._store.query(page_query) 

149 if not events: 

150 break 

151 

152 buf = io.StringIO() 

153 writer = csv.writer(buf) 

154 for event in events: 

155 writer.writerow(_event_to_csv_row(event)) 

156 yield buf.getvalue().encode() 

157 

158 if len(events) < _EXPORT_PAGE_SIZE: 

159 break 

160 offset += _EXPORT_PAGE_SIZE 

161 

162 logger.debug("audit_export_csv_complete", offset=offset) 

163 

164 async def export_json(self, query: AuditQuery) -> AsyncIterator[bytes]: 

165 """Stream audit records as newline-delimited JSON (NDJSON). 

166 

167 Each line is a complete JSON object representing one 

168 :class:`~lexigram.ai.governance.audit.models.AIAuditEvent`. 

169 

170 Args: 

171 query: Filter criteria; ``limit`` and ``offset`` are used for 

172 internal pagination. 

173 

174 Yields: 

175 UTF-8 encoded bytes for each page of NDJSON lines. 

176 """ 

177 offset = 0 

178 page_query = AuditQuery( 

179 start=query.start, 

180 end=query.end, 

181 event_types=query.event_types, 

182 user_id=query.user_id, 

183 model=query.model, 

184 provider=query.provider, 

185 status=query.status, 

186 limit=_EXPORT_PAGE_SIZE, 

187 offset=offset, 

188 ) 

189 

190 while True: 

191 page_query.offset = offset 

192 events = await self._store.query(page_query) 

193 if not events: 

194 break 

195 

196 lines: list[str] = [] 

197 for event in events: 

198 record = { 

199 "event_id": event.event_id, 

200 "event_type": event.event_type.value, 

201 "timestamp": event.timestamp.isoformat(), 

202 "model": event.model, 

203 "provider": event.provider, 

204 "user_id": event.user_id, 

205 "status": event.status, 

206 "tokens": event.tokens, 

207 "cost": event.cost, 

208 "latency_ms": event.latency_ms, 

209 "metadata": event.metadata, 

210 } 

211 lines.append(dumps_str(record)) 

212 yield ("\n".join(lines) + "\n").encode() 

213 

214 if len(events) < _EXPORT_PAGE_SIZE: 

215 break 

216 offset += _EXPORT_PAGE_SIZE 

217 

218 logger.debug("audit_export_json_complete", offset=offset) 

219 

220 async def summary(self, window_hours: int = 24) -> AuditSummary: 

221 """Aggregate audit statistics for the last *window_hours* hours. 

222 

223 Args: 

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

225 

226 Returns: 

227 :class:`~lexigram.ai.governance.audit.models.AuditSummary` with 

228 total events, spend, tokens, denied count, and breakdowns by 

229 model, user, and event type. 

230 """ 

231 start = datetime.now(UTC) - timedelta(hours=window_hours) 

232 q = AuditQuery(start=start, limit=0) 

233 logger.debug("audit_summary", window_hours=window_hours) 

234 return await self._store.aggregate(q) 

235 

236 

237__all__ = ["AuditQueryService"]