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

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 

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

14 

15from __future__ import annotations 

16 

17import asyncio 

18from collections.abc import AsyncGenerator, AsyncIterator, Mapping 

19import time 

20from typing import Any, Literal, cast 

21 

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 

67 

68__all__ = ["RelayGatewayService"] 

69 

70logger = get_logger(__name__) 

71 

72 

73class RelayGatewayService: 

74 """Relay request lifecycle (buffered and streaming). 

75 

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. 

79 

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

95 

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. 

110 

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 

136 

137 async def handle( 

138 self, request: RelayGatewayRequest 

139 ) -> Result[RelayGatewayResult, RelayGatewayError]: 

140 """Run the buffered or streaming relay lifecycle for one request. 

141 

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. 

148 

149 Args: 

150 request: The gateway request. 

151 

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 

196 

197 async def _dispatch( 

198 self, request: RelayGatewayRequest 

199 ) -> tuple[Result[RelayGatewayResult, RelayGatewayError], str]: 

200 """Run the ordered dependency pipeline for one request. 

201 

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 

422 

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. 

427 

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. 

432 

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 ) 

546 

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. 

558 

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. 

565 

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) 

657 

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. 

667 

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. 

674 

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 ) 

689 

690 @staticmethod 

691 def _usage_from_snapshot(snapshot: object) -> RelayUsage | None: 

692 """Extract ``RelayUsage`` from a session snapshot when present. 

693 

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. 

698 

699 Args: 

700 snapshot: The session snapshot returned by 

701 ``RelayStreamSessionProtocol.snapshot``. 

702 

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 ) 

720 

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. 

728 

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`. 

734 

735 Args: 

736 request: The gateway request being dispatched. 

737 billing: The billing lifecycle to reserve through. 

738 channel: The selected channel. 

739 

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()) 

768 

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. 

778 

779 Settlement failures are logged and never propagate: the response 

780 path has already completed by the time accounting runs. 

781 

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 ) 

799 

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. 

805 

806 Args: 

807 request: The gateway request being dispatched. 

808 channel: The selected channel. 

809 

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 ) 

821 

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. 

829 

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 ) 

842 

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. 

854 

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 ) 

875 

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. 

884 

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. 

890 

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 

925 

926 def _upstream_url(self, channel: RelayChannel, model: str) -> str: 

927 """Build the endpoint URL for *channel*'s target format. 

928 

929 Args: 

930 channel: The selected channel. 

931 model: Outbound model alias (embedded in the Gemini path). 

932 

933 Returns: 

934 The standard endpoint path for the channel's target format 

935 joined onto the channel's base URL. 

936 

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

952 

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 )