Coverage for src / lexigram / ai / relay / gateway / operations / metrics.py: 100%
57 statements
« prev ^ index » next coverage.py v7.13.5, created at 2026-08-08 23:08 +0800
« prev ^ index » next coverage.py v7.13.5, created at 2026-08-08 23:08 +0800
1"""Route metrics aggregation for the relay gateway operations surface.
3``RelayMetricsService`` turns routed operational events into stable
4``RelayRouteMetrics`` rows grouped by directed route and time window,
5counting conversion losses from stable codes. Registry diagnostics and
6package-version reporting stay import-free for mapper modules.
7"""
9from __future__ import annotations
11from collections.abc import Mapping, Sequence
12from dataclasses import dataclass
13from datetime import datetime
14from typing import Literal, Protocol, runtime_checkable
16from lexigram.contracts.ai.relay import (
17 ConversionQuality,
18 RelayFormat,
19 RelayGatewayError,
20 RelayRegistryDiagnostics,
21 RelayRegistryProtocol,
22 RelayRouteMetrics,
23 TimeWindow,
24)
26__all__ = [
27 "CONVERTER_ID",
28 "RelayMetricsService",
29 "RelayRouteEvent",
30 "RelayRouteEventSourceProtocol",
31 "RouteMetricEventKind",
32]
34CONVERTER_ID = "relay-converter"
35"""Diagnostics identifier for the built-in relay converter."""
37RouteMetricEventKind = Literal[
38 "request_completed",
39 "conversion_loss",
40 "unsupported_feature",
41 "stream_cancelled",
42 "stream_timeout",
43 "stream_truncated",
44]
47@dataclass(frozen=True, slots=True)
48class RelayRouteEvent:
49 """One operational event feeding route metric aggregation.
51 Attributes:
52 kind: Event kind; only ``conversion_loss`` events carry a code.
53 source: Source wire format of the route.
54 target: Target wire format of the route.
55 occurred_at: When the event happened (UTC).
56 loss_code: Stable conversion loss code for ``conversion_loss``
57 events; ignored for every other kind.
58 """
60 kind: RouteMetricEventKind
61 source: RelayFormat
62 target: RelayFormat
63 occurred_at: datetime
64 loss_code: str | None = None
67@runtime_checkable
68class RelayRouteEventSourceProtocol(Protocol):
69 """Provides operational events inside a bounded window."""
71 async def events(self, window: TimeWindow) -> Sequence[RelayRouteEvent]:
72 """Return the events observed inside *window*.
74 Args:
75 window: Bounded aggregation window.
77 Returns:
78 Events bounded to the window; an empty sequence means no
79 activity, never a failed lookup.
80 """
81 ...
84class RelayMetricsService:
85 """Aggregate route metrics for the admin operations surface.
87 Counts, per directed route and window:
89 - ``request_count`` from ``request_completed`` events.
90 - ``loss_counts`` from ``conversion_loss`` codes (never free-form
91 messages).
92 - ``unsupported_count`` from ``unsupported_feature`` events.
93 - ``stream_failure_count`` from stream
94 cancelled/timeout/truncated events.
96 A missing event source is a failed dependency; an empty source
97 yields a stable empty result, never a fabricated zero-count row.
98 """
100 def __init__(
101 self,
102 events: RelayRouteEventSourceProtocol | None,
103 converter: RelayRegistryProtocol | None = None,
104 registration_errors: tuple[str, ...] = (),
105 ) -> None:
106 """Bind the metrics service to its dependencies.
108 Args:
109 events: Operational event source for route aggregation.
110 ``None`` makes ``route_metrics`` a failed dependency.
111 converter: Optional converter registry used for route-quality
112 and diagnostics. ``None`` makes
113 ``registry_diagnostics`` a failed dependency.
114 registration_errors: Registrations that failed at wiring
115 time, surfaced verbatim in diagnostics.
116 """
117 self._events = events
118 self._converter = converter
119 self._registration_errors = registration_errors
121 async def route_metrics(self, window: TimeWindow) -> Sequence[RelayRouteMetrics]:
122 """Return per-route metrics aggregated inside *window*.
124 Args:
125 window: Bounded aggregation window.
127 Returns:
128 One row per route that saw activity within the window.
130 Raises:
131 RelayGatewayError: With ``DEPENDENCY_UNAVAILABLE`` when no
132 event source is registered.
133 """
134 if self._events is None:
135 raise RelayGatewayError(
136 code="DEPENDENCY_UNAVAILABLE",
137 message="relay route event source is not registered",
138 status_code=503,
139 request_id="",
140 )
141 events = await self._events.events(window)
142 buckets: dict[tuple[RelayFormat, RelayFormat], list[RelayRouteEvent]] = {}
143 for event in events:
144 if window.start < event.occurred_at < window.end:
145 key = (event.source, event.target)
146 buckets.setdefault(key, []).append(event)
147 rows: list[RelayRouteMetrics] = []
148 for (source, target), route_events in sorted(buckets.items()):
149 rows.append(
150 RelayRouteMetrics(
151 source=source,
152 target=target,
153 quality=self._quality(source, target),
154 request_count=sum(
155 1 for e in route_events if e.kind == "request_completed"
156 ),
157 loss_counts=self._losses(route_events),
158 unsupported_count=sum(
159 1 for e in route_events if e.kind == "unsupported_feature"
160 ),
161 stream_failure_count=sum(
162 1
163 for e in route_events
164 if e.kind
165 in ("stream_cancelled", "stream_timeout", "stream_truncated")
166 ),
167 converter_id=self._converter_id(source, target),
168 window_start=window.start,
169 window_end=window.end,
170 )
171 )
172 return rows
174 async def registry_diagnostics(self) -> RelayRegistryDiagnostics:
175 """Return converter capability diagnostics.
177 Returns:
178 Converter identifier, version, mapper ids, supported route
179 pairs, and startup registration failures.
181 Raises:
182 RelayGatewayError: With ``DEPENDENCY_UNAVAILABLE`` when no
183 converter registry is registered.
184 """
185 if self._converter is None:
186 raise RelayGatewayError(
187 code="DEPENDENCY_UNAVAILABLE",
188 message="converter registry is not registered",
189 status_code=503,
190 request_id="",
191 )
192 return RelayRegistryDiagnostics(
193 converter_id=CONVERTER_ID,
194 converter_version=self._converter.converter_version(),
195 mapper_ids=self._converter.mapper_ids(),
196 supported_routes=self._converter.converter_routes(),
197 registration_errors=self._registration_errors,
198 )
200 @staticmethod
201 def _losses(
202 route_events: Sequence[RelayRouteEvent],
203 ) -> Mapping[str, int]:
204 """Count conversion losses by stable code."""
205 counts: dict[str, int] = {}
206 for event in route_events:
207 if event.kind == "conversion_loss" and event.loss_code:
208 counts[event.loss_code] = counts.get(event.loss_code, 0) + 1
209 return counts
211 def _quality(self, source: RelayFormat, target: RelayFormat) -> ConversionQuality:
212 """Return the route quality from the converter matrix."""
213 if self._converter is None:
214 return ConversionQuality.DISCOURAGED
215 return self._converter.route_quality(source, target)
217 def _converter_id(self, source: RelayFormat, target: RelayFormat) -> str | None:
218 """Return the route converter identifier, when known."""
219 if self._converter is None:
220 return None
221 return f"{source.value}_to_{target.value}"