Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-relay-gateway/src/lexigram/ai/relay/gateway/service.py: 42%

55 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

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