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

1"""Route metrics aggregation for the relay gateway operations surface. 

2 

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

8 

9from __future__ import annotations 

10 

11from collections.abc import Mapping, Sequence 

12from dataclasses import dataclass 

13from datetime import datetime 

14from typing import Literal, Protocol, runtime_checkable 

15 

16from lexigram.contracts.ai.relay import ( 

17 ConversionQuality, 

18 RelayFormat, 

19 RelayGatewayError, 

20 RelayRegistryDiagnostics, 

21 RelayRegistryProtocol, 

22 RelayRouteMetrics, 

23 TimeWindow, 

24) 

25 

26__all__ = [ 

27 "CONVERTER_ID", 

28 "RelayMetricsService", 

29 "RelayRouteEvent", 

30 "RelayRouteEventSourceProtocol", 

31 "RouteMetricEventKind", 

32] 

33 

34CONVERTER_ID = "relay-converter" 

35"""Diagnostics identifier for the built-in relay converter.""" 

36 

37RouteMetricEventKind = Literal[ 

38 "request_completed", 

39 "conversion_loss", 

40 "unsupported_feature", 

41 "stream_cancelled", 

42 "stream_timeout", 

43 "stream_truncated", 

44] 

45 

46 

47@dataclass(frozen=True, slots=True) 

48class RelayRouteEvent: 

49 """One operational event feeding route metric aggregation. 

50 

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

59 

60 kind: RouteMetricEventKind 

61 source: RelayFormat 

62 target: RelayFormat 

63 occurred_at: datetime 

64 loss_code: str | None = None 

65 

66 

67@runtime_checkable 

68class RelayRouteEventSourceProtocol(Protocol): 

69 """Provides operational events inside a bounded window.""" 

70 

71 async def events(self, window: TimeWindow) -> Sequence[RelayRouteEvent]: 

72 """Return the events observed inside *window*. 

73 

74 Args: 

75 window: Bounded aggregation window. 

76 

77 Returns: 

78 Events bounded to the window; an empty sequence means no 

79 activity, never a failed lookup. 

80 """ 

81 ... 

82 

83 

84class RelayMetricsService: 

85 """Aggregate route metrics for the admin operations surface. 

86 

87 Counts, per directed route and window: 

88 

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. 

95 

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

99 

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. 

107 

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 

120 

121 async def route_metrics(self, window: TimeWindow) -> Sequence[RelayRouteMetrics]: 

122 """Return per-route metrics aggregated inside *window*. 

123 

124 Args: 

125 window: Bounded aggregation window. 

126 

127 Returns: 

128 One row per route that saw activity within the window. 

129 

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 

173 

174 async def registry_diagnostics(self) -> RelayRegistryDiagnostics: 

175 """Return converter capability diagnostics. 

176 

177 Returns: 

178 Converter identifier, version, mapper ids, supported route 

179 pairs, and startup registration failures. 

180 

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 ) 

199 

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 

210 

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) 

216 

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