1"""Relay gateway request lifecycle (buffered and streaming).
2
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.
9
10The pipeline's specialized concerns live in focused modules:
11
12- :mod:`lexigram.ai.relay.gateway.validation` — request boundary checks.
13- :mod:`lexigram.ai.relay.gateway.buffered` — buffered dispatch pipeline.
14- :mod:`lexigram.ai.relay.gateway.streaming` — streaming preflight and
15 lazy stream consumption with exactly-once settlement.
16- :mod:`lexigram.ai.relay.gateway.operations.billing` — admission and
17 settlement lifecycle.
18- :mod:`lexigram.ai.relay.gateway.operations.upstream` — endpoint URLs,
19 model resolution, upstream calls, and failover accounting.
20- :mod:`lexigram.ai.relay.gateway.operations.telemetry` — structured
21 request events.
22
23``config.provider_options`` are intentionally not merged into
24``RelayOptions`` yet because no mapping schema exists.
25"""
26
27from __future__ import annotations
28
29import time
30
31from lexigram.ai.relay.gateway.buffered import BufferedDispatchMixin
32from lexigram.ai.relay.gateway.channels import RelayChannelRegistry
33from lexigram.ai.relay.gateway.codec import RelayPayloadCodec
34from lexigram.ai.relay.gateway.config import RelayGatewayConfig
35from lexigram.ai.relay.gateway.errors import unexpected_error
36from lexigram.ai.relay.gateway.operations import telemetry
37from lexigram.ai.relay.gateway.operations.failover import RelayFailoverTracker
38from lexigram.ai.relay.gateway.operations.streams import RelayStreamRegistry
39from lexigram.ai.relay.gateway.streaming import StreamingMixin
40from lexigram.ai.relay.gateway.upstream import HTTPUpstreamAdapter
41from lexigram.ai.relay.gateway.validation import validate_gateway_request
42from lexigram.contracts.ai.governance import RelayBillingProtocol
43from lexigram.contracts.ai.relay import (
44 MediaResolverProtocol,
45 RelayConverterProtocol,
46 RelayGatewayError,
47 RelayGatewayRequest,
48 RelayGatewayResult,
49)
50from lexigram.contracts.auth.guard import AuthorizerProtocol
51from lexigram.contracts.core.result import Err, Result
52from lexigram.logging import get_logger
53
54__all__ = ["RelayGatewayService", "validate_gateway_request"]
55
56logger = get_logger(__name__)
57
58
59class RelayGatewayService(BufferedDispatchMixin, StreamingMixin):
60 """Relay request lifecycle (buffered and streaming).
61
62 The service is stateless between requests and never touches request
63 headers, payloads, or upstream details in error messages; errors are
64 always safe ``RelayGatewayError`` values.
65
66 Attributes:
67 _converter: Engine implementing ``RelayConverterProtocol``.
68 _codec: Wire DTO codec.
69 _registry: Deterministic channel selector.
70 _upstream: HTTP transport adapter.
71 _config: Gateway configuration (channel table and model suffixes).
72 _authorizer: Optional authorization check before dispatch.
73 _billing: Optional billing lifecycle; when ``None`` admission and
74 settlement are skipped.
75 _media_resolver: Optional URL-media resolver threaded into the
76 conversion context.
77 _streams: Optional registry of active streams used to expose
78 in-flight streams and cancel handles to operators; ``None``
79 disables stream registration (but not streaming itself).
80 _failover: Optional consecutive-failure tracker; when ``None``
81 upstream failures never affect runtime selection state.
82 """
83
84 def __init__(
85 self,
86 converter: RelayConverterProtocol,
87 codec: RelayPayloadCodec,
88 registry: RelayChannelRegistry,
89 upstream: HTTPUpstreamAdapter,
90 config: RelayGatewayConfig,
91 *,
92 authorizer: AuthorizerProtocol | None = None,
93 billing: RelayBillingProtocol | None = None,
94 media_resolver: MediaResolverProtocol | None = None,
95 streams: RelayStreamRegistry | None = None,
96 failover: RelayFailoverTracker | None = None,
97 ) -> None:
98 """Bind the service to its dependencies.
99
100 Args:
101 converter: Conversion engine for request/response payloads.
102 codec: Wire DTO decoder/encoder.
103 registry: Channel selection registry.
104 upstream: Upstream transport adapter.
105 config: Static gateway configuration.
106 authorizer: Optional authorizer; when ``None`` authorization
107 is skipped.
108 billing: Optional billing lifecycle; when ``None`` the
109 gateway runs without admission control or settlement.
110 media_resolver: Optional URL-media resolver placed on the
111 conversion context; ``None`` disables media resolution.
112 streams: Optional stream registry for operator visibility
113 and forced cancellation; ``None`` keeps streaming
114 functional without registry bookkeeping.
115 failover: Optional consecutive-failure tracker; when ``None``
116 upstream failures never affect runtime selection state.
117 """
118 self._converter = converter
119 self._codec = codec
120 self._registry = registry
121 self._upstream = upstream
122 self._config = config
123 self._authorizer = authorizer
124 self._billing = billing
125 self._media_resolver = media_resolver
126 self._streams = streams
127 self._failover = failover
128
129 async def handle(
130 self, request: RelayGatewayRequest
131 ) -> Result[RelayGatewayResult, RelayGatewayError]:
132 """Run the buffered or streaming relay lifecycle for one request.
133
134 Dependencies run in fixed order: authorize, select channel,
135 reserve billing capacity, convert request, then either call
136 upstream, decode, convert response back, and settle (buffered),
137 or create the stream session and hand back a lazy stream that
138 consumes upstream events and settles when exhausted (streaming).
139 Any preflight failure short-circuits the pipeline.
140
141 Args:
142 request: The gateway request.
143
144 Returns:
145 ``Ok(RelayGatewayResult)`` on success, or
146 ``Err(RelayGatewayError)`` on the first failure. Unexpected
147 exceptions from dependencies never escape: they are logged
148 and mapped to a generic ``CONVERSION_FAILED`` error.
149 """
150 started = time.monotonic()
151 validation_error = validate_gateway_request(request)
152 if validation_error is not None:
153 logger.warning(
154 "relay_gateway_invalid_request",
155 request_id=request.request_id,
156 error=validation_error.message,
157 )
158 telemetry.log_request_completed(request, "", validation_error, started)
159 return Err(validation_error)
160 logger.info(
161 "relay_gateway_request_accepted",
162 request_id=request.request_id,
163 tenant_id=request.tenant_id,
164 source=request.source,
165 model=request.model,
166 stream=request.stream,
167 )
168 try:
169 if request.stream:
170 result, channel_name = await self._handle_streaming(request)
171 else:
172 result, channel_name = await self._dispatch(request)
173 except Exception as exc:
174 logger.warning(
175 "relay_gateway_unexpected_error",
176 request_id=request.request_id,
177 error=str(exc),
178 )
179 error = unexpected_error(request.request_id)
180 telemetry.log_request_completed(request, "", error, started)
181 return Err(error)
182 if result.is_err():
183 telemetry.log_request_completed(
184 request, channel_name, result.unwrap_err(), started
185 )
186 return result
187 outcome = result.unwrap()
188 telemetry.log_request_completed(
189 request,
190 channel_name,
191 outcome,
192 started,
193 target=outcome.metadata.target if outcome.metadata else "",
194 loss_codes=outcome.metadata.loss_codes if outcome.metadata else (),
195 )
196 return result