Coverage for src / lexigram / ai / relay / engine.py: 86%
102 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"""Public relay conversion engine implementing :class:`RelayConverterProtocol`.
3The engine is synchronous, side-effect free, and performs no HTTP,
4channel selection, billing, or model selection. Callers supply the
5already-selected upstream model and any host callbacks via
6:class:`RelayConversionContext`.
7"""
9from __future__ import annotations
11from typing import Any
13from lexigram.ai.relay.context import ConversionContext
14from lexigram.ai.relay.errors import (
15 translate,
16 unsupported_feature,
17 unsupported_format,
18 unsupported_route,
19)
20from lexigram.ai.relay.mappers.base import warning_messages
21from lexigram.ai.relay.quality import route_quality
22from lexigram.ai.relay.registry import RelayConverterRegistry, Route, RouteSpec
23from lexigram.contracts.ai.exceptions import RelayError
24from lexigram.contracts.ai.relay.context import RelayConversionContext
25from lexigram.contracts.ai.relay.protocols import (
26 RelayConverterProtocol,
27 RelayRegistryProtocol,
28 RelayStreamOptions,
29 RelayStreamSessionProtocol,
30)
31from lexigram.contracts.ai.relay.types import (
32 ConversionQuality,
33 RelayConvertResult,
34 RelayFormat,
35 RelayRequestPayload,
36 RelayResponsePayload,
37)
38from lexigram.contracts.core.result import Err, Ok, Result
40__all__ = [
41 "RelayConverterEngine",
42 "convert_request_by_id",
43 "convert_request_via",
44 "convert_response_by_id",
45 "convert_response_via",
46]
49def _converter_id(source: RelayFormat, target: RelayFormat) -> str:
50 """Return the stable ``"<source>_to_<target>"`` converter identifier."""
51 return f"{source.value}_to_{target.value}"
54def _normalize_steps(source: RelayFormat, target: RelayFormat) -> tuple[str, ...]:
55 """Return the canonical conversion path: source, IR, target."""
56 return (source.value, "canonical_ir", target.value)
59def _route_spec(mapper: Any) -> RouteSpec | None:
60 """Extract route metadata from a registry-produced mapper, or ``None``."""
61 spec = getattr(mapper, "spec", None)
62 return spec if isinstance(spec, RouteSpec) else None
65class RelayConverterEngine(RelayConverterProtocol):
66 """Explicit source/target conversion over a caller-owned registry.
68 Args:
69 registry: The registry used when a conversion call does not pass
70 its own ``registry`` override.
71 """
73 def __init__(self, registry: RelayRegistryProtocol) -> None:
74 """Initialise the engine with a default registry."""
75 self.registry: RelayRegistryProtocol = registry
77 def convert_request(
78 self,
79 payload: RelayRequestPayload,
80 source: RelayFormat,
81 target: RelayFormat,
82 *,
83 context: RelayConversionContext | None = None,
84 registry: RelayRegistryProtocol | None = None,
85 ) -> Result[RelayConvertResult[RelayRequestPayload], RelayError]:
86 """Convert a request payload from *source* to *target*."""
87 conversion_context = ConversionContext.wrap(context)
88 try:
89 return self._convert(
90 payload, source, target, registry, conversion_context, response=False
91 )
92 except Exception as exc: # noqa: BLE001
93 return Err(translate(exc, detail="convert_request"))
95 def convert_response(
96 self,
97 payload: RelayResponsePayload,
98 source: RelayFormat,
99 target: RelayFormat,
100 *,
101 context: RelayConversionContext | None = None,
102 registry: RelayRegistryProtocol | None = None,
103 ) -> Result[RelayConvertResult[RelayResponsePayload], RelayError]:
104 """Convert a non-stream response payload from *source* to *target*."""
105 conversion_context = ConversionContext.wrap(context)
106 try:
107 return self._convert(
108 payload, source, target, registry, conversion_context, response=True
109 )
110 except Exception as exc: # noqa: BLE001
111 return Err(translate(exc, detail="convert_response"))
113 def new_stream_session(
114 self,
115 source: RelayFormat,
116 target: RelayFormat,
117 *,
118 options: RelayStreamOptions | None = None,
119 context: RelayConversionContext | None = None,
120 registry: RelayRegistryProtocol | None = None,
121 ) -> Result[RelayStreamSessionProtocol, RelayError]:
122 """Create a stateful stream session for one upstream stream.
124 Stateful stream conversion is delivered by the shared stream
125 lifecycle in a later task; every current route reports no stream
126 support.
127 """
128 return Err(
129 unsupported_feature("stateful stream conversion is not implemented yet")
130 )
132 def convert_stream_chunk(
133 self,
134 session: RelayStreamSessionProtocol,
135 event: Any,
136 ) -> tuple[Any, ...]:
137 """Convert one source stream event through *session*.
139 Raises:
140 RelayError: Stateful stream conversion is not implemented yet.
141 """
142 raise unsupported_feature("stateful stream conversion is not implemented yet")
144 def finalize(
145 self,
146 session: RelayStreamSessionProtocol,
147 ) -> tuple[Any, ...]:
148 """Close a stream deterministically and return terminal events.
150 Raises:
151 RelayError: Stateful stream conversion is not implemented yet.
152 """
153 raise unsupported_feature("stateful stream conversion is not implemented yet")
155 def _convert(
156 self,
157 payload: Any,
158 source: RelayFormat,
159 target: RelayFormat,
160 registry: RelayRegistryProtocol | None,
161 context: ConversionContext,
162 *,
163 response: bool,
164 ) -> Result[RelayConvertResult[Any], RelayError]:
165 """Shared request/response conversion path."""
166 if not isinstance(source, RelayFormat):
167 return Err(
168 unsupported_format(f"source must be a RelayFormat, got {source!r}")
169 )
170 if not isinstance(target, RelayFormat):
171 return Err(
172 unsupported_format(f"target must be a RelayFormat, got {target!r}")
173 )
174 if source is target:
175 return Ok(self._noop_result(payload, source, target, context))
177 active_registry = registry or self.registry
178 route_mapper = active_registry.mapper(source, target)
179 if route_mapper is None:
180 return Err(
181 unsupported_route(
182 f"no relay route from {source.value} to {target.value}"
183 )
184 )
185 route = route_mapper if isinstance(route_mapper, Route) else None
186 spec = route.spec if route is not None else _route_spec(route_mapper)
187 if spec is not None:
188 supported = spec.response_supported if response else spec.request_supported
189 if not supported:
190 return Err(
191 unsupported_feature(
192 f"{'response' if response else 'request'} conversion "
193 f"{source.value} -> {target.value} is not supported"
194 )
195 )
196 if response:
197 if route is not None:
198 ir = route.response_to_ir(payload, context=context)
199 else:
200 ir = route_mapper.response_to_ir(payload)
201 if ir.is_err():
202 return ir
203 source_ir = ir.unwrap()
204 if route is not None:
205 target_payload = route.ir_to_response(source_ir, context=context)
206 else:
207 target_payload = route_mapper.ir_to_response(source_ir)
208 if target_payload.is_err():
209 return target_payload
210 usage = source_ir.usage
211 else:
212 if route is not None:
213 ir = route.request_to_ir(payload, context=context)
214 else:
215 ir = route_mapper.request_to_ir(payload)
216 if ir.is_err():
217 return ir
218 source_ir = ir.unwrap()
219 if route is not None:
220 target_payload = route.ir_to_request(source_ir, context=context)
221 else:
222 target_payload = route_mapper.ir_to_request(source_ir)
223 if target_payload.is_err():
224 return target_payload
225 usage = None
226 losses = tuple(context.losses)
227 return Ok(
228 RelayConvertResult(
229 value=target_payload.unwrap(),
230 source=source,
231 target=target,
232 converter_id=_converter_id(source, target),
233 quality=route_quality(source, target),
234 steps=_normalize_steps(source, target),
235 usage=usage,
236 losses=losses,
237 warnings=warning_messages(losses),
238 )
239 )
241 def _noop_result(
242 self,
243 value: Any,
244 source: RelayFormat,
245 target: RelayFormat,
246 context: ConversionContext,
247 ) -> RelayConvertResult[Any]:
248 """Return the original typed payload with a no-op result."""
249 losses = tuple(context.losses)
250 return RelayConvertResult(
251 value=value,
252 source=source,
253 target=target,
254 converter_id=_converter_id(source, target),
255 quality=ConversionQuality.GOOD,
256 steps=(),
257 usage=None,
258 losses=losses,
259 warnings=warning_messages(losses),
260 )
263# ---------------------------------------------------------------------------
264# Explicit-path helpers
265# ---------------------------------------------------------------------------
268def convert_request_via(
269 registry: RelayRegistryProtocol,
270 payload: RelayRequestPayload,
271 source: RelayFormat,
272 target: RelayFormat,
273 *,
274 context: RelayConversionContext | None = None,
275) -> Result[RelayConvertResult[RelayRequestPayload], RelayError]:
276 """Convert a request through an explicit registry.
278 Args:
279 registry: The registry to resolve the route from.
280 payload: The source request DTO.
281 source: Source wire format.
282 target: Target wire format.
283 context: Optional host conversion context.
285 Returns:
286 A request conversion result, or a relay error.
287 """
288 return RelayConverterEngine(registry).convert_request(
289 payload, source, target, context=context
290 )
293def convert_response_via(
294 registry: RelayRegistryProtocol,
295 payload: RelayResponsePayload,
296 source: RelayFormat,
297 target: RelayFormat,
298 *,
299 context: RelayConversionContext | None = None,
300) -> Result[RelayConvertResult[RelayResponsePayload], RelayError]:
301 """Convert a response through an explicit registry.
303 Args:
304 registry: The registry to resolve the route from.
305 payload: The source response DTO.
306 source: Source wire format.
307 target: Target wire format.
308 context: Optional host conversion context.
310 Returns:
311 A response conversion result, or a relay error.
312 """
313 return RelayConverterEngine(registry).convert_response(
314 payload, source, target, context=context
315 )
318def convert_request_by_id(
319 registry: RelayConverterRegistry,
320 payload: RelayRequestPayload,
321 converter_id: str,
322 *,
323 context: RelayConversionContext | None = None,
324) -> Result[RelayConvertResult[RelayRequestPayload], RelayError]:
325 """Convert a request through the route identified by ``converter_id``.
327 Args:
328 registry: The registry carrying the route.
329 payload: The source request DTO.
330 converter_id: A stable ``"<source>_to_<target>"`` identifier.
331 context: Optional host conversion context.
333 Returns:
334 A request conversion result, or a relay error.
335 """
336 spec = registry.route_by_id(converter_id)
337 if spec is None:
338 return Err(unsupported_route(f"no relay route with id {converter_id!r}"))
339 return RelayConverterEngine(registry).convert_request(
340 payload, spec.source, spec.target, context=context
341 )
344def convert_response_by_id(
345 registry: RelayConverterRegistry,
346 payload: RelayResponsePayload,
347 converter_id: str,
348 *,
349 context: RelayConversionContext | None = None,
350) -> Result[RelayConvertResult[RelayResponsePayload], RelayError]:
351 """Convert a response through the route identified by ``converter_id``.
353 Args:
354 registry: The registry carrying the route.
355 payload: The source response DTO.
356 converter_id: A stable ``"<source>_to_<target>"`` identifier.
357 context: Optional host conversion context.
359 Returns:
360 A response conversion result, or a relay error.
361 """
362 spec = registry.route_by_id(converter_id)
363 if spec is None:
364 return Err(unsupported_route(f"no relay route with id {converter_id!r}"))
365 return RelayConverterEngine(registry).convert_response(
366 payload, spec.source, spec.target, context=context
367 )