Coverage for src / lexigram / ai / relay / gateway / service.py: 94%
262 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"""Relay gateway request lifecycle (buffered and streaming).
3``RelayGatewayService`` orchestrates one request through the
4dependencies: authorization, channel selection, billing admission, request
5conversion, the protected upstream call, response conversion, billing
6settlement, and result metadata assembly. Streaming requests run the
7same preflight and then consume the upstream SSE stream lazily through
8the stream session, settling billing exactly once when the stream ends.
10Provider authentication headers are out of scope until Task 7.
11``config.provider_options`` are intentionally not merged into
12``RelayOptions`` yet because no mapping schema exists.
13"""
15from __future__ import annotations
17import asyncio
18from collections.abc import AsyncGenerator, AsyncIterator, Mapping
19import time
20from typing import Any, Literal, cast
22from lexigram.ai.relay.gateway.channels import RelayChannelRegistry
23from lexigram.ai.relay.gateway.codec import RelayPayloadCodec
24from lexigram.ai.relay.gateway.config import RelayGatewayConfig
25from lexigram.ai.relay.gateway.errors import (
26 auth_denied,
27 billing_error_to_gateway,
28 conversion_error_to_gateway,
29 with_request_id,
30)
31from lexigram.ai.relay.gateway.operations.streams import RelayStreamRegistry
32from lexigram.ai.relay.gateway.stream import UpstreamEventParser, relay_stream
33from lexigram.ai.relay.gateway.upstream import HTTPUpstreamAdapter
34from lexigram.contracts.ai.exceptions import RelayError
35from lexigram.contracts.ai.governance import (
36 RelayBillingProtocol,
37 RelayUsageReservation,
38 RelayUsageScope,
39)
40from lexigram.contracts.ai.relay import (
41 ConversionQuality,
42 MediaResolverProtocol,
43 RelayChannel,
44 RelayConversionContext,
45 RelayConverterProtocol,
46 RelayConvertResult,
47 RelayFormat,
48 RelayGatewayError,
49 RelayGatewayMetadata,
50 RelayGatewayRequest,
51 RelayGatewayResult,
52 RelayLoss,
53 RelayOptions,
54 RelayRequestPayload,
55 RelayStreamSessionProtocol,
56 RelayUpstreamProtocol,
57 RelayUsage,
58 RelayWireEvent,
59 UpstreamRequest,
60 UpstreamResponse,
61)
62from lexigram.contracts.ai.relay.gateway import RelayGatewayErrorCode
63from lexigram.contracts.auth.guard import AuthorizerProtocol
64from lexigram.contracts.core.result import Err, Ok, Result
65from lexigram.logging import get_logger
66from lexigram.serialization import dumps
68__all__ = ["RelayGatewayService"]
70logger = get_logger(__name__)
73class RelayGatewayService:
74 """Relay request lifecycle (buffered and streaming).
76 The service is stateless between requests and never touches request
77 headers, payloads, or upstream details in error messages; errors are
78 always safe ``RelayGatewayError`` values.
80 Attributes:
81 _converter: Engine implementing ``RelayConverterProtocol``.
82 _codec: Wire DTO codec.
83 _registry: Deterministic channel selector.
84 _upstream: HTTP transport adapter.
85 _config: Gateway configuration (channel table and model suffixes).
86 _authorizer: Optional authorization check before dispatch.
87 _billing: Optional billing lifecycle; when ``None`` admission and
88 settlement are skipped.
89 _media_resolver: Optional URL-media resolver threaded into the
90 conversion context.
91 _streams: Optional registry of active streams used to expose
92 in-flight streams and cancel handles to operators; ``None``
93 disables stream registration (but not streaming itself).
94 """
96 def __init__(
97 self,
98 converter: RelayConverterProtocol,
99 codec: RelayPayloadCodec,
100 registry: RelayChannelRegistry,
101 upstream: HTTPUpstreamAdapter,
102 config: RelayGatewayConfig,
103 *,
104 authorizer: AuthorizerProtocol | None = None,
105 billing: RelayBillingProtocol | None = None,
106 media_resolver: MediaResolverProtocol | None = None,
107 streams: RelayStreamRegistry | None = None,
108 ) -> None:
109 """Bind the service to its dependencies.
111 Args:
112 converter: Conversion engine for request/response payloads.
113 codec: Wire DTO decoder/encoder.
114 registry: Channel selection registry.
115 upstream: Upstream transport adapter.
116 config: Static gateway configuration.
117 authorizer: Optional authorizer; when ``None`` authorization
118 is skipped.
119 billing: Optional billing lifecycle; when ``None`` the
120 gateway runs without admission control or settlement.
121 media_resolver: Optional URL-media resolver placed on the
122 conversion context; ``None`` disables media resolution.
123 streams: Optional stream registry for operator visibility
124 and forced cancellation; ``None`` keeps streaming
125 functional without registry bookkeeping.
126 """
127 self._converter = converter
128 self._codec = codec
129 self._registry = registry
130 self._upstream = upstream
131 self._config = config
132 self._authorizer = authorizer
133 self._billing = billing
134 self._media_resolver = media_resolver
135 self._streams = streams
137 async def handle(
138 self, request: RelayGatewayRequest
139 ) -> Result[RelayGatewayResult, RelayGatewayError]:
140 """Run the buffered or streaming relay lifecycle for one request.
142 Dependencies run in fixed order: authorize, select channel,
143 reserve billing capacity, convert request, then either call
144 upstream, decode, convert response back, and settle (buffered),
145 or create the stream session and hand back a lazy stream that
146 consumes upstream events and settles when exhausted (streaming).
147 Any preflight failure short-circuits the pipeline.
149 Args:
150 request: The gateway request.
152 Returns:
153 ``Ok(RelayGatewayResult)`` on success, or
154 ``Err(RelayGatewayError)`` on the first failure. Unexpected
155 exceptions from dependencies never escape: they are logged
156 and mapped to a generic ``CONVERSION_FAILED`` error.
157 """
158 started = time.monotonic()
159 logger.info(
160 "relay_gateway_request_accepted",
161 request_id=request.request_id,
162 tenant_id=request.tenant_id,
163 source=request.source,
164 model=request.model,
165 stream=request.stream,
166 )
167 try:
168 if request.stream:
169 result, channel_name = await self._handle_streaming(request)
170 else:
171 result, channel_name = await self._dispatch(request)
172 except Exception as exc:
173 logger.warning(
174 "relay_gateway_unexpected_error",
175 request_id=request.request_id,
176 error=str(exc),
177 )
178 error = self._unexpected_error(request.request_id)
179 self._log_request_completed(request, "", error, started)
180 return Err(error)
181 if result.is_err():
182 self._log_request_completed(
183 request, channel_name, result.unwrap_err(), started
184 )
185 return result
186 outcome = result.unwrap()
187 self._log_request_completed(
188 request,
189 channel_name,
190 outcome,
191 started,
192 target=outcome.metadata.target if outcome.metadata else "",
193 loss_codes=outcome.metadata.loss_codes if outcome.metadata else (),
194 )
195 return result
197 async def _dispatch(
198 self, request: RelayGatewayRequest
199 ) -> tuple[Result[RelayGatewayResult, RelayGatewayError], str]:
200 """Run the ordered dependency pipeline for one request.
202 Returns:
203 ``tuple`` of the pipeline result and the selected channel
204 name. The channel name is ``""`` when selection failed
205 before a channel was chosen.
206 """
207 if self._authorizer is not None:
208 allowed = await self._authorizer.authorize(
209 user=request.tenant_id,
210 action="relay.invoke",
211 resource=request.model,
212 )
213 if not allowed:
214 return Err(auth_denied(request.request_id)), ""
215 max_attempts = self._config.max_upstream_retries + 1
216 tried: set[str] = set()
217 last_upstream_error: RelayGatewayError | None = None
218 last_channel_name = ""
219 last_channel: RelayChannel | None = None
220 held_reservation: RelayUsageReservation | None = None
221 for attempt in range(1, max_attempts + 1):
222 selected = self._registry.select(
223 source=request.source,
224 model=request.model,
225 stream=request.stream,
226 preferred=request.channel.name if request.channel else None,
227 exclude=frozenset(tried),
228 )
229 if selected.is_err():
230 if last_upstream_error is not None:
231 billing = self._billing
232 if (
233 held_reservation is not None
234 and billing is not None
235 and last_channel is not None
236 ):
237 await self._settle(
238 billing,
239 held_reservation,
240 self._empty_settle_result(request, last_channel),
241 status="failed",
242 )
243 return (
244 Err(with_request_id(last_upstream_error, request.request_id)),
245 last_channel_name,
246 )
247 return (
248 Err(with_request_id(selected.unwrap_err(), request.request_id)),
249 "",
250 )
251 channel = selected.unwrap()
252 last_channel = channel
253 last_channel_name = channel.name
254 logger.info(
255 "relay_gateway_channel_selected",
256 request_id=request.request_id,
257 channel=channel.name,
258 target_format=channel.target_format,
259 model=request.model,
260 )
261 billing = self._billing
262 reservation: RelayUsageReservation | None = None
263 if billing is not None:
264 admitted = await self._pre_consume(request, billing, channel)
265 if admitted.is_err():
266 return Err(admitted.unwrap_err()), channel.name
267 reservation = admitted.unwrap()
268 if held_reservation is not None:
269 await billing.release(held_reservation)
270 held_reservation = None
271 outbound_model = request.model + self._config.model_suffix.get(
272 channel.name, ""
273 )
274 context = RelayConversionContext(
275 request_id=request.request_id,
276 channel_name=channel.name,
277 upstream_model=outbound_model,
278 options=RelayOptions(),
279 media_resolver=self._media_resolver,
280 )
281 conv = self._converter.convert_request(
282 payload=cast("RelayRequestPayload", request.payload),
283 source=request.source,
284 target=channel.target_format,
285 context=context,
286 )
287 if conv.is_err():
288 if reservation is not None and billing is not None:
289 await billing.release(reservation)
290 return (
291 Err(
292 conversion_error_to_gateway(
293 conv.unwrap_err(), request.request_id
294 )
295 ),
296 channel.name,
297 )
298 converted_request = conv.unwrap()
299 self._log_conversion_loss(
300 request.request_id,
301 converted_request.converter_id,
302 converted_request.losses,
303 )
304 upstream_response = await self._call_upstream(
305 channel, outbound_model, converted_request.value.to_dict(), request
306 )
307 if upstream_response.is_err():
308 upstream_error = upstream_response.unwrap_err()
309 if upstream_error.retryable and attempt < max_attempts:
310 tried.add(channel.name)
311 last_upstream_error = upstream_error
312 held_reservation = reservation
313 logger.info(
314 "relay_gateway_upstream_retry",
315 request_id=request.request_id,
316 channel=channel.name,
317 error_code=upstream_error.code,
318 attempt=attempt,
319 )
320 continue
321 if reservation is not None and billing is not None:
322 await self._settle(
323 billing,
324 reservation,
325 self._empty_settle_result(request, channel),
326 status="failed",
327 )
328 return (
329 Err(with_request_id(upstream_error, request.request_id)),
330 channel.name,
331 )
332 resp = upstream_response.unwrap()
333 if resp.payload is None:
334 if reservation is not None and billing is not None:
335 await self._settle(
336 billing,
337 reservation,
338 self._empty_settle_result(request, channel),
339 status="failed",
340 )
341 return (
342 Err(
343 RelayGatewayError(
344 code=RelayGatewayErrorCode.UPSTREAM_MALFORMED,
345 message="malformed upstream response",
346 status_code=502,
347 request_id=request.request_id,
348 retryable=False,
349 )
350 ),
351 channel.name,
352 )
353 decoded = self._codec.decode_response_payload(
354 target=channel.target_format,
355 data=dict(resp.payload),
356 request_id=request.request_id,
357 )
358 if decoded.is_err():
359 if reservation is not None and billing is not None:
360 await self._settle(
361 billing,
362 reservation,
363 self._empty_settle_result(request, channel),
364 status="failed",
365 )
366 return Err(decoded.unwrap_err()), channel.name
367 back = self._converter.convert_response(
368 payload=decoded.unwrap(),
369 source=channel.target_format,
370 target=request.source,
371 context=context,
372 )
373 if back.is_err():
374 if reservation is not None and billing is not None:
375 await self._settle(
376 billing,
377 reservation,
378 self._empty_settle_result(request, channel),
379 status="failed",
380 )
381 return (
382 Err(
383 conversion_error_to_gateway(
384 back.unwrap_err(), request.request_id
385 )
386 ),
387 channel.name,
388 )
389 converted = back.unwrap()
390 self._log_conversion_loss(
391 request.request_id, converted.converter_id, converted.losses
392 )
393 if reservation is not None and billing is not None:
394 await self._settle(billing, reservation, converted, status="completed")
395 metadata = RelayGatewayMetadata(
396 converter_id=converted.converter_id,
397 source=request.source,
398 target=channel.target_format,
399 quality=converted.quality,
400 loss_codes=tuple(loss.reason for loss in converted.losses),
401 warnings=converted.warnings,
402 )
403 return (
404 Ok(
405 RelayGatewayResult(
406 status_code=resp.status_code,
407 headers={**resp.headers, "x-request-id": request.request_id},
408 payload=converted.value.to_dict(),
409 stream=None,
410 metadata=metadata,
411 )
412 ),
413 channel.name,
414 )
415 error = last_upstream_error or RelayGatewayError(
416 code=RelayGatewayErrorCode.CHANNEL_DISABLED,
417 message="no channels available",
418 status_code=404,
419 request_id=request.request_id,
420 )
421 return Err(with_request_id(error, request.request_id)), last_channel_name
423 async def _handle_streaming(
424 self, request: RelayGatewayRequest
425 ) -> tuple[Result[RelayGatewayResult, RelayGatewayError], str]:
426 """Run the streaming preflight and return a lazy stream result.
428 Authorization, channel selection, billing admission, request
429 conversion, and stream-session creation all complete before the
430 first frame is delivered; upstream I/O happens lazily as the
431 returned stream is consumed by the caller.
433 Returns:
434 ``tuple`` of the pipeline result (whose ``stream`` holds the
435 lazy ``AsyncIterator`` on success) and the selected channel
436 name. The channel name is ``""`` when selection failed
437 before a channel was chosen.
438 """
439 if self._authorizer is not None:
440 allowed = await self._authorizer.authorize(
441 user=request.tenant_id,
442 action="relay.invoke",
443 resource=request.model,
444 )
445 if not allowed:
446 return Err(auth_denied(request.request_id)), ""
447 selected = self._registry.select(
448 source=request.source,
449 model=request.model,
450 stream=True,
451 preferred=request.channel.name if request.channel else None,
452 )
453 if selected.is_err():
454 return (
455 Err(with_request_id(selected.unwrap_err(), request.request_id)),
456 "",
457 )
458 channel = selected.unwrap()
459 logger.info(
460 "relay_gateway_channel_selected",
461 request_id=request.request_id,
462 channel=channel.name,
463 target_format=channel.target_format,
464 model=request.model,
465 )
466 billing = self._billing
467 reservation: RelayUsageReservation | None = None
468 if billing is not None:
469 admitted = await self._pre_consume(request, billing, channel)
470 if admitted.is_err():
471 return Err(admitted.unwrap_err()), channel.name
472 reservation = admitted.unwrap()
473 outbound_model = request.model + self._config.model_suffix.get(channel.name, "")
474 context = RelayConversionContext(
475 request_id=request.request_id,
476 channel_name=channel.name,
477 upstream_model=outbound_model,
478 options=RelayOptions(),
479 media_resolver=self._media_resolver,
480 )
481 conv = self._converter.convert_request(
482 payload=cast("RelayRequestPayload", request.payload),
483 source=request.source,
484 target=channel.target_format,
485 context=context,
486 )
487 if conv.is_err():
488 if reservation is not None and billing is not None:
489 await billing.release(reservation)
490 return (
491 Err(conversion_error_to_gateway(conv.unwrap_err(), request.request_id)),
492 channel.name,
493 )
494 converted_request = conv.unwrap()
495 session = self._converter.new_stream_session(
496 source=request.source,
497 target=channel.target_format,
498 context=context,
499 )
500 if session.is_err():
501 if reservation is not None and billing is not None:
502 await billing.release(reservation)
503 return (
504 Err(
505 conversion_error_to_gateway(
506 session.unwrap_err(), request.request_id
507 )
508 ),
509 channel.name,
510 )
511 stream_session = session.unwrap()
512 self._log_conversion_loss(
513 request.request_id,
514 converted_request.converter_id,
515 converted_request.losses,
516 )
517 metadata = RelayGatewayMetadata(
518 converter_id=converted_request.converter_id,
519 source=request.source,
520 target=channel.target_format,
521 quality=converted_request.quality,
522 loss_codes=tuple(loss.reason for loss in converted_request.losses),
523 warnings=converted_request.warnings,
524 )
525 stream = self._stream_events(
526 request,
527 channel,
528 outbound_model,
529 converted_request,
530 reservation,
531 stream_session,
532 context,
533 )
534 return (
535 Ok(
536 RelayGatewayResult(
537 status_code=200,
538 headers={"x-request-id": request.request_id},
539 payload=None,
540 stream=stream,
541 metadata=metadata,
542 )
543 ),
544 channel.name,
545 )
547 async def _stream_events(
548 self,
549 request: RelayGatewayRequest,
550 channel: RelayChannel,
551 outbound_model: str,
552 converted_request: RelayConvertResult[RelayRequestPayload],
553 reservation: RelayUsageReservation | None,
554 session: RelayStreamSessionProtocol,
555 context: RelayConversionContext,
556 ) -> AsyncIterator[RelayWireEvent]:
557 """Consume the upstream stream and settle the reservation once.
559 Each consumer pull forwards exactly one upstream chunk through
560 the session; cancellation, truncation, and malformed framing
561 follow the ``relay_stream`` lifecycle. The ``finally`` block
562 runs when the consumer ends the stream (completion, disconnect,
563 or error) and settles billing exactly once from the session
564 snapshot.
566 Yields:
567 Normalized ``RelayWireEvent`` values framed by the stream
568 session.
569 """
570 streams = self._streams
571 stream_id: str | None = None
572 cancel_handle: asyncio.Event | None = None
573 if streams is not None:
574 stream_id, cancel_handle = streams.register(
575 channel=channel.name,
576 model=outbound_model,
577 request_id=request.request_id,
578 )
579 parser = UpstreamEventParser(
580 session=session,
581 source=channel.target_format,
582 request_id=request.request_id,
583 )
584 url = self._upstream_url(channel, outbound_model)
585 payload = (
586 converted_request.value.to_dict()
587 if converted_request.value is not None
588 else {}
589 )
590 logger.info(
591 "relay_gateway_stream_started",
592 request_id=request.request_id,
593 channel=channel.name,
594 method="POST",
595 url=url,
596 )
597 upstream_request = UpstreamRequest(
598 request_id=request.request_id,
599 method="POST",
600 url=url,
601 headers={"content-type": "application/json"},
602 payload=dict(payload),
603 timeout_seconds=channel.timeout_seconds,
604 channel_name=channel.name,
605 )
606 truncated = False
607 stream_iter: AsyncGenerator[RelayWireEvent, None] | None = None
608 try:
609 stream_iter = cast(
610 "AsyncGenerator[RelayWireEvent, None]",
611 relay_stream(
612 cast(
613 "RelayUpstreamProtocol",
614 self._upstream,
615 ),
616 upstream_request,
617 parser,
618 cancel_handle=cancel_handle,
619 ),
620 )
621 try:
622 async for wire in stream_iter:
623 yield wire
624 except (RelayGatewayError, RelayError) as error:
625 logger.warning(
626 "relay_gateway_stream_malformed",
627 request_id=request.request_id,
628 channel=channel.name,
629 error=str(error),
630 )
631 truncated = True
632 raise
633 finally:
634 if stream_iter is not None:
635 await stream_iter.aclose()
636 if truncated:
637 status: Literal["completed", "failed", "cancelled", "truncated"] = (
638 "truncated"
639 )
640 elif parser.cancelled:
641 status = "cancelled"
642 elif parser.truncated:
643 status = "truncated"
644 else:
645 status = "completed"
646 if streams is not None and stream_id is not None:
647 streams.unregister(stream_id)
648 billing = self._billing
649 if billing is not None and reservation is not None:
650 settled = self._stream_settle_result(
651 request,
652 channel,
653 converter_id=converted_request.converter_id,
654 session=session,
655 )
656 await self._settle(billing, reservation, settled, status=status)
658 def _stream_settle_result(
659 self,
660 request: RelayGatewayRequest,
661 channel: RelayChannel,
662 *,
663 converter_id: str,
664 session: RelayStreamSessionProtocol,
665 ) -> RelayConvertResult[Any]:
666 """Build the settled result from the stream session snapshot.
668 Args:
669 request: The gateway request being settled.
670 channel: The selected channel.
671 converter_id: Converter that produced the stream session.
672 session: The stream session whose snapshot carries the
673 settled usage.
675 Returns:
676 A ``RelayConvertResult`` carrying normalized usage extracted
677 from the session snapshot (or no usage when the snapshot
678 exposes none).
679 """
680 usage = self._usage_from_snapshot(session.snapshot())
681 return RelayConvertResult(
682 value=None,
683 source=request.source,
684 target=channel.target_format,
685 converter_id=converter_id,
686 quality=ConversionQuality.GOOD,
687 usage=usage,
688 )
690 @staticmethod
691 def _usage_from_snapshot(snapshot: object) -> RelayUsage | None:
692 """Extract ``RelayUsage`` from a session snapshot when present.
694 Snapshots are opaque; only a ``Mapping`` carrying a ``usage``
695 sub-mapping is inspected, accepting either the OpenAI-style
696 ``prompt_tokens``/``completion_tokens`` keys or the
697 Claude-style ``input_tokens``/``output_tokens`` keys.
699 Args:
700 snapshot: The session snapshot returned by
701 ``RelayStreamSessionProtocol.snapshot``.
703 Returns:
704 A normalized ``RelayUsage`` when the snapshot exposes one,
705 else ``None`` (settlement then records no usage).
706 """
707 if not isinstance(snapshot, Mapping):
708 return None
709 usage = snapshot.get("usage")
710 if not isinstance(usage, Mapping):
711 return None
712 prompt = usage.get("prompt_tokens", usage.get("input_tokens"))
713 completion = usage.get("completion_tokens", usage.get("output_tokens"))
714 if not isinstance(prompt, int) or not isinstance(completion, int):
715 return None
716 return RelayUsage(
717 prompt_tokens=prompt,
718 completion_tokens=completion,
719 )
721 async def _pre_consume(
722 self,
723 request: RelayGatewayRequest,
724 billing: RelayBillingProtocol,
725 channel: RelayChannel,
726 ) -> Result[RelayUsageReservation, RelayGatewayError]:
727 """Reserve billing capacity before conversion and upstream I/O.
729 The inbound payload is re-decoded into a typed request DTO so the
730 billing pipeline can estimate prompt and output budgets; a
731 payload that rejects decoding fails the request here, before any
732 upstream I/O. Billing denials short-circuit the pipeline and are
733 classified through :func:`billing_error_to_gateway`.
735 Args:
736 request: The gateway request being dispatched.
737 billing: The billing lifecycle to reserve through.
738 channel: The selected channel.
740 Returns:
741 ``Ok(reservation)`` when admission is proven, or
742 ``Err(RelayGatewayError)`` carrying the classified failure.
743 """
744 scope = RelayUsageScope(
745 tenant_id=request.tenant_id,
746 model=request.model,
747 channel=channel.name,
748 )
749 dto = self._codec.decode_request(
750 source=request.source,
751 raw=dumps(dict(request.payload)),
752 request_id=request.request_id,
753 )
754 if dto.is_err():
755 return Err(dto.unwrap_err())
756 admitted = await billing.pre_consume(request.request_id, scope, dto.unwrap())
757 if admitted.is_err():
758 error = admitted.unwrap_err()
759 logger.warning(
760 "relay_gateway_billing_denied",
761 request_id=request.request_id,
762 channel=channel.name,
763 code=error.code,
764 error=error.message,
765 )
766 return Err(billing_error_to_gateway(error, request.request_id))
767 return Ok(admitted.unwrap())
769 async def _settle(
770 self,
771 billing: RelayBillingProtocol,
772 reservation: RelayUsageReservation,
773 result: RelayConvertResult[Any],
774 *,
775 status: Literal["completed", "failed", "cancelled", "truncated"],
776 ) -> None:
777 """Settle the reservation exactly once without failing the response.
779 Settlement failures are logged and never propagate: the response
780 path has already completed by the time accounting runs.
782 Args:
783 billing: The billing lifecycle to settle through.
784 reservation: The reservation granted by ``pre_consume``.
785 result: The conversion result carrying settled usage, or an
786 empty result when the attempt produced no billable usage.
787 status: Terminal lifecycle status of the attempt.
788 """
789 settled = await billing.settle(reservation, result, status=status)
790 if settled.is_err():
791 error = settled.unwrap_err()
792 logger.warning(
793 "relay_gateway_settle_failed",
794 request_id=reservation.request_id,
795 status=status,
796 code=error.code,
797 error=error.message,
798 )
800 @staticmethod
801 def _empty_settle_result(
802 request: RelayGatewayRequest, channel: RelayChannel
803 ) -> RelayConvertResult[Any]:
804 """Build the usage-free result settled for failed attempts.
806 Args:
807 request: The gateway request being dispatched.
808 channel: The selected channel.
810 Returns:
811 A ``RelayConvertResult`` carrying no usage so the billing
812 pipeline records an attempted-but-unbilled attempt.
813 """
814 return RelayConvertResult(
815 value=None,
816 source=request.source,
817 target=channel.target_format,
818 converter_id="",
819 quality=ConversionQuality.GOOD,
820 )
822 def _log_conversion_loss(
823 self,
824 request_id: str,
825 converter_id: str,
826 losses: tuple[RelayLoss, ...] | list[RelayLoss],
827 ) -> None:
828 """Emit the conversion-loss event when a conversion recorded losses.
830 Args:
831 request_id: The gateway request identifier.
832 converter_id: The converter that produced the losses.
833 losses: Recorded conversion losses.
834 """
835 if losses:
836 logger.info(
837 "relay_gateway_conversion_loss",
838 request_id=request_id,
839 converter_id=converter_id,
840 loss_codes=tuple(loss.reason for loss in losses),
841 )
843 def _log_request_completed(
844 self,
845 request: RelayGatewayRequest,
846 channel_name: str,
847 outcome: RelayGatewayResult | RelayGatewayError,
848 started: float,
849 *,
850 target: str = "",
851 loss_codes: tuple[str, ...] = (),
852 ) -> None:
853 """Emit the terminal request-completed event for any outcome.
855 Args:
856 request: The original gateway request.
857 channel_name: Selected channel name (or ``""`` when unknown).
858 outcome: The success result or the error that ended the flow.
859 started: Monotonic start time used to compute the duration.
860 target: Target format name (success path only).
861 loss_codes: Conversion loss codes (success path only).
862 """
863 logger.info(
864 "relay_gateway_request_completed",
865 request_id=request.request_id,
866 tenant_id=request.tenant_id,
867 channel=channel_name,
868 source=request.source,
869 target=target,
870 status_code=outcome.status_code,
871 code=outcome.code if isinstance(outcome, RelayGatewayError) else "OK",
872 duration_ms=round((time.monotonic() - started) * 1000, 2),
873 loss_codes=loss_codes,
874 )
876 async def _call_upstream(
877 self,
878 channel: RelayChannel,
879 outbound_model: str,
880 payload: dict[str, Any],
881 request: RelayGatewayRequest,
882 ) -> Result[UpstreamResponse, RelayGatewayError]:
883 """Send the converted payload to the selected channel's endpoint.
885 Args:
886 channel: The selected channel.
887 outbound_model: Model alias with the channel's suffix applied.
888 payload: Converted request payload dict.
889 request: The original gateway request.
891 Returns:
892 ``Ok(UpstreamResponse)`` or ``Err`` as returned by the
893 adapter; the adapter already normalizes transport failures.
894 """
895 url = self._upstream_url(channel, outbound_model)
896 logger.info(
897 "relay_gateway_upstream_started",
898 request_id=request.request_id,
899 channel=channel.name,
900 method="POST",
901 url=url,
902 )
903 upstream = await self._upstream.request(
904 UpstreamRequest(
905 request_id=request.request_id,
906 method="POST",
907 url=url,
908 headers={"content-type": "application/json"},
909 payload=payload,
910 timeout_seconds=channel.timeout_seconds,
911 channel_name=channel.name,
912 )
913 )
914 if upstream.is_err():
915 err = upstream.unwrap_err()
916 logger.warning(
917 "relay_gateway_upstream_failed",
918 request_id=request.request_id,
919 channel=channel.name,
920 code=err.code,
921 status_code=err.status_code,
922 error=str(err),
923 )
924 return upstream
926 def _upstream_url(self, channel: RelayChannel, model: str) -> str:
927 """Build the endpoint URL for *channel*'s target format.
929 Args:
930 channel: The selected channel.
931 model: Outbound model alias (embedded in the Gemini path).
933 Returns:
934 The standard endpoint path for the channel's target format
935 joined onto the channel's base URL.
937 Raises:
938 ValueError: The channel's target format is not one of the
939 four relay wire formats. Unreachable via registry
940 validation.
941 """
942 base = channel.upstream_base_url.rstrip("/")
943 if channel.target_format == RelayFormat.OPENAI_CHAT:
944 return f"{base}/v1/chat/completions"
945 if channel.target_format == RelayFormat.OPENAI_RESPONSES:
946 return f"{base}/v1/responses"
947 if channel.target_format == RelayFormat.CLAUDE:
948 return f"{base}/v1/messages"
949 if channel.target_format == RelayFormat.GEMINI:
950 return f"{base}/v1beta/models/{model}:generateContent"
951 raise ValueError(f"unsupported target format: {channel.target_format}")
953 @staticmethod
954 def _unexpected_error(request_id: str) -> RelayGatewayError:
955 """Build the generic error for unexpected dependency failures."""
956 return RelayGatewayError(
957 code=RelayGatewayErrorCode.CONVERSION_FAILED,
958 message="Unexpected relay gateway failure",
959 status_code=500,
960 request_id=request_id,
961 retryable=False,
962 )