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