Skip to content
Open
Show file tree
Hide file tree
Changes from 9 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ POSTGRES_PYTEST_TARGETS := \
tests/integration/test_repositories.py::test_accounts_upsert_with_merge_enabled_serializes_concurrent_same_email \
tests/integration/test_sticky_sessions_api.py::test_durable_bridge_owned_alias_registration_is_epoch_fenced \
tests/integration/test_proxy_api_extended.py::test_proxy_stream_usage_limit_returns_http_error \
tests/integration/test_api_keys_api.py::test_rate_limit_header_failure_releases_reservation_once \
tests/integration/test_codex_usage_api.py::test_codex_usage_aggregates_windows \
tests/integration/test_proxy_compact.py::test_proxy_compact_headers_include_monthly_only_credits \
tests/integration/test_repositories.py::test_accounts_upsert_with_merge_disabled_uses_identity_lock_on_postgresql \
Expand Down
9 changes: 7 additions & 2 deletions app/modules/proxy/_service/http_bridge/account_sessions.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
from __future__ import annotations

from app.modules.proxy._service.http_bridge.helpers import _extract_model_class, _log_http_bridge_event
from app.modules.proxy._service.http_bridge.helpers import (
_claim_http_bridge_session_close,
_extract_model_class,
_log_http_bridge_event,
)
from app.modules.proxy._service.http_bridge.protocol import _HTTPBridgeServiceProtocol
from app.modules.proxy._service.support import _HTTPBridgeSession

Expand All @@ -26,5 +30,6 @@ async def close_http_bridge_sessions_for_account(self: _HTTPBridgeServiceProtoco
sessions_to_close.append(detached)

for session in sessions_to_close:
await self._close_http_bridge_session_bounded(session, reason="account_binding_changed")
if _claim_http_bridge_session_close(session):
await self._close_http_bridge_session_bounded(session, reason="account_binding_changed")
return len(sessions_to_close)
21 changes: 16 additions & 5 deletions app/modules/proxy/_service/http_bridge/helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -675,6 +675,13 @@ def _http_bridge_session_has_visible_requests(session: "_HTTPBridgeSession") ->
)


def _claim_http_bridge_session_close(session: "_HTTPBridgeSession") -> bool:
if session.upstream_close_attempted:
return False
session.upstream_close_attempted = True
return True


async def _close_http_bridge_session_bounded(
service: Any,
session: "_HTTPBridgeSession",
Expand All @@ -687,14 +694,18 @@ async def _close_http_bridge_session_bounded(
service._close_http_bridge_session(session),
name=f"http-bridge-close-{_hash_identifier(session.key.affinity_key)}",
)
service._background_cleanup_tasks.add(close_task)

def untrack_close(done_task: asyncio.Task[None]) -> None:
service._background_cleanup_tasks.discard(done_task)

close_task.add_done_callback(untrack_close)

def track_after_interruption(*, interruption: str) -> None:
def observe_after_interruption(*, interruption: str) -> None:
if close_task.done():
return
service._background_cleanup_tasks.add(close_task)

def close_done(done_task: asyncio.Task[None]) -> None:
service._background_cleanup_tasks.discard(done_task)
try:
done_task.result()
except asyncio.CancelledError:
Expand Down Expand Up @@ -729,7 +740,7 @@ def close_done(done_task: asyncio.Task[None]) -> None:
timeout=_HTTP_BRIDGE_BACKGROUND_CLOSE_TIMEOUT_SECONDS,
)
except TimeoutError:
track_after_interruption(interruption="timeout")
observe_after_interruption(interruption="timeout")
logger.warning(
"http_bridge_session_close_timeout reason=%s bridge_kind=%s bridge_key=%s "
"account_id=%s model=%s timeout_seconds=%.1f background_cleanup_tasks=%d",
Expand All @@ -742,7 +753,7 @@ def close_done(done_task: asyncio.Task[None]) -> None:
len(service._background_cleanup_tasks),
)
except asyncio.CancelledError:
track_after_interruption(interruption="cancellation")
observe_after_interruption(interruption="cancellation")
raise
except Exception:
logger.warning(
Expand Down
44 changes: 34 additions & 10 deletions app/modules/proxy/_service/http_bridge/mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@
_HTTP_BRIDGE_INFLIGHT_STARTED_AT_ATTR,
_active_http_bridge_instance_ring,
_await_task_deferring_cancellation,
_claim_http_bridge_session_close,
_close_http_bridge_session_bounded,
_durable_bridge_lookup_active_owner,
_durable_bridge_lookup_allows_local_reuse,
Expand Down Expand Up @@ -252,6 +253,8 @@ def _schedule_http_bridge_session_closes(
reason: str,
) -> None:
for session in sessions:
if not _claim_http_bridge_session_close(session):
continue
if len(self._background_cleanup_tasks) >= _HTTP_BRIDGE_BACKGROUND_CLEANUP_WARN_THRESHOLD:
logger.warning(
"http_bridge_background_cleanup_backlog action=session_close count=%d threshold=%d reason=%s",
Expand Down Expand Up @@ -1320,7 +1323,8 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None:
owns_creation = True
try:
for session_to_close in sessions_to_close_before_create:
await self._close_http_bridge_session_bounded(session_to_close, reason="registry_detach")
if _claim_http_bridge_session_close(session_to_close):
await self._close_http_bridge_session_bounded(session_to_close, reason="registry_detach")
except BaseException as exc:
if owns_creation:
await self._fail_http_bridge_inflight_session_creation(key, inflight_future, exc)
Expand Down Expand Up @@ -1536,7 +1540,11 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None:
inflight_future.set_exception(exc)
inflight_future.exception()
if created_session is not None and not session_registered:
await self._close_http_bridge_session(created_session)
if _claim_http_bridge_session_close(created_session):
await self._close_http_bridge_session_bounded(
created_session,
reason="registration_failed",
)
raise
assert created_session is not None
_log_http_bridge_event(
Expand Down Expand Up @@ -1583,14 +1591,30 @@ async def close_all_http_bridge_sessions(self) -> None:
error_type="server_error",
),
)
for inflight_future in inflight_futures:
if inflight_future.done():
continue
inflight_future.set_exception(shutdown_error)
inflight_future.exception()
for session in sessions_to_close:
await self._close_http_bridge_session(session)
await self._drain_http_bridge_background_cleanup_tasks(reason="shutdown")
owned_sessions = [session for session in sessions_to_close if _claim_http_bridge_session_close(session)]
Comment thread
mastertyko marked this conversation as resolved.
Outdated

async def close_owned_sessions() -> None:
try:
for session in owned_sessions:
await self._close_http_bridge_session_bounded(session, reason="shutdown")
finally:
await self._drain_http_bridge_background_cleanup_tasks(reason="shutdown")

cleanup_task = asyncio.create_task(
close_owned_sessions(),
name="http-bridge-cleanup-shutdown",
)
cancellation: asyncio.CancelledError | None = None
try:
for inflight_future in inflight_futures:
if inflight_future.done():
continue
inflight_future.set_exception(shutdown_error)
inflight_future.exception()
finally:
_, cancellation = await _await_task_deferring_cancellation(cleanup_task)
if cancellation is not None:
raise cancellation

async def mark_http_bridge_draining(self) -> None:
try:
Expand Down
58 changes: 37 additions & 21 deletions app/modules/proxy/_service/http_bridge/request_submit.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@
from app.modules.proxy._service.http_bridge.helpers import (
_await_task_deferring_cancellation,
_build_http_bridge_prewarm_text,
_claim_http_bridge_session_close,
_http_bridge_key_strength,
_http_bridge_precreated_retry_failure_error,
_http_bridge_prewarm_enabled,
Expand Down Expand Up @@ -1412,7 +1413,7 @@ async def _retire_http_bridge_after_drain_if_ready(self: Any, session: "_HTTPBri
)
if should_reconnect:
session.pending_requests.clear()
session.upstream_close_attempted = True
should_reconnect = _claim_http_bridge_session_close(session)
if not should_reconnect:
return False

Expand All @@ -1426,27 +1427,42 @@ async def _retire_stale_pending_http_bridge_session(
detail: str,
) -> None:
session.closed = True
async with self._http_bridge_lock:
if self._http_bridge_sessions.get(session.key) is session:
self._http_bridge_sessions.pop(session.key, None)
self._unregister_http_bridge_turn_states_locked(session)
self._unregister_http_bridge_previous_response_ids_locked(session)
async with session.pending_lock:
should_close = not session.upstream_close_attempted
if should_close:
session.upstream_close_attempted = True
if should_close:
await self._close_http_bridge_session_bounded(session, reason="retire_stale_pending")
_log_http_bridge_event(
"retire_stale_pending",
session.key,
account_id=session.account.id,
model=session.request_model,
pending_count=await self._http_bridge_pending_count(session),
detail=detail,
cache_key_family=session.key.affinity_kind,
model_class=_extract_model_class(session.request_model) if session.request_model else None,
if session.upstream_reader is asyncio.current_task():
session.upstream_reader = None
owns_close = _claim_http_bridge_session_close(session)

async def retire_session() -> None:
async with self._http_bridge_lock:
if self._http_bridge_sessions.get(session.key) is session:
self._http_bridge_sessions.pop(session.key, None)
self._unregister_http_bridge_turn_states_locked(session)
self._unregister_http_bridge_previous_response_ids_locked(session)
if owns_close:
await self._close_http_bridge_session_bounded(session, reason="retire_stale_pending")
_log_http_bridge_event(
"retire_stale_pending",
session.key,
account_id=session.account.id,
model=session.request_model,
pending_count=self._http_bridge_pending_count_nowait(session, context="retire_stale_pending_log"),
detail=detail,
cache_key_family=session.key.affinity_kind,
model_class=_extract_model_class(session.request_model) if session.request_model else None,
)

retire_task = asyncio.create_task(
retire_session(),
name=f"http-bridge-close-retire-{_hash_identifier(session.key.affinity_key)}",
)
self._background_cleanup_tasks.add(retire_task)

def untrack_retirement(done_task: asyncio.Task[None]) -> None:
self._background_cleanup_tasks.discard(done_task)

retire_task.add_done_callback(untrack_retirement)
_, cancellation = await _await_task_deferring_cancellation(retire_task)
if cancellation is not None:
raise cancellation

async def _retry_http_bridge_request_on_fresh_upstream(
self: Any,
Expand Down
42 changes: 29 additions & 13 deletions app/modules/proxy/_service/http_bridge/streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,8 @@
_sticky_key_from_compact_payload as _sticky_key_from_compact_payload,
)
from app.modules.proxy._service.http_bridge.helpers import (
_await_task_deferring_cancellation,
_claim_http_bridge_session_close,
_effective_http_bridge_idle_ttl_seconds,
_http_bridge_durable_lookup_allows_turn_state_takeover,
_http_bridge_is_context_overflow_error,
Expand Down Expand Up @@ -2442,22 +2444,36 @@ async def _reset_http_bridge_session_after_local_terminal_error(
error_code: str,
error_message: str,
) -> None:
owns_close = False
async with self._http_bridge_lock:
if self._http_bridge_sessions.get(session.key) is session:
self._http_bridge_sessions.pop(session.key, None)
async with session.pending_lock:
session.queued_request_count = 0
await self._fail_pending_websocket_requests(
account=session.account,
account_id_value=session.account.id,
pending_requests=session.pending_requests,
pending_lock=session.pending_lock,
error_code=error_code,
error_message=error_message,
api_key=None,
response_create_gate=session.response_create_gate,
)
await self._close_http_bridge_session(session)
owns_close = _claim_http_bridge_session_close(session)
close_cancellation: asyncio.CancelledError | None = None
try:
async with session.pending_lock:
session.queued_request_count = 0
await self._fail_pending_websocket_requests(
account=session.account,
account_id_value=session.account.id,
pending_requests=session.pending_requests,
pending_lock=session.pending_lock,
error_code=error_code,
error_message=error_message,
api_key=None,
response_create_gate=session.response_create_gate,
)
finally:
if owns_close:
close_task = asyncio.create_task(
self._close_http_bridge_session_bounded(
session,
reason="local_terminal_error",
)
)
_, close_cancellation = await _await_task_deferring_cancellation(close_task)
if close_cancellation is not None:
raise close_cancellation

async def _stream_http_bridge_session_events(
self: Any,
Expand Down
50 changes: 46 additions & 4 deletions app/modules/proxy/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
from typing import Any, Final, Literal, Protocol, cast
from uuid import uuid4

import anyio
from fastapi import (
APIRouter,
Body,
Expand Down Expand Up @@ -1958,6 +1959,39 @@ async def _rate_limit_headers_for_request(
return await context.service.rate_limit_headers()


async def _release_reservation_deferring_cancellation(
reservation: ApiKeyUsageReservationData,
) -> None:
with anyio.CancelScope(shield=True):
task = asyncio.create_task(_release_reservation(reservation))
while True:
try:
await asyncio.shield(task)
return
except asyncio.CancelledError:
if task.cancelled():
raise


async def _rate_limit_headers_with_reservation_cleanup(
context: ProxyContext,
api_key: ApiKeyData | None,
owned_reservation: ApiKeyUsageReservationData | None,
) -> dict[str, str]:
try:
return await _rate_limit_headers_for_request(context, api_key)
except BaseException:
if owned_reservation is not None:
try:
await _release_reservation_deferring_cancellation(owned_reservation)
except (Exception, asyncio.CancelledError):
logger.warning(
"Failed to release API key reservation after rate-limit header failure",
exc_info=True,
)
raise


def _select_codex_usage_limit(
limits: list[V1UsageLimitResponse],
window: str,
Expand Down Expand Up @@ -4875,7 +4909,15 @@ async def _stream_responses(
)
)

rate_limit_headers = await _rate_limit_headers_for_request(context, api_key) if include_rate_limit_headers else {}
rate_limit_headers = (
await _rate_limit_headers_with_reservation_cleanup(
context,
api_key,
reservation if owns_reservation else None,
)
if include_rate_limit_headers
else {}
)
bridge_active = prefer_http_bridge and proxy_service_module.get_settings().http_responses_session_bridge_enabled
effective_headers = forwarded_headers or request.headers
client_ip = forwarded_client_ip if forwarded_request else resolve_request_client_host(request)
Expand Down Expand Up @@ -5080,7 +5122,7 @@ async def _collect_responses(
request_usage_budget=estimate_api_key_request_usage(payload),
)

rate_limit_headers = await _rate_limit_headers_for_request(context, api_key)
rate_limit_headers = await _rate_limit_headers_with_reservation_cleanup(context, api_key, reservation)
bridge_active = prefer_http_bridge and proxy_service_module.get_settings().http_responses_session_bridge_enabled
downstream_turn_state = (
proxy_affinity_module.ensure_http_downstream_turn_state(request.headers) if bridge_active else None
Expand Down Expand Up @@ -5240,7 +5282,7 @@ async def _compact_responses(
request_usage_budget=request_usage_budget,
)

rate_limit_headers = await _rate_limit_headers_for_request(context, api_key)
rate_limit_headers = await _rate_limit_headers_with_reservation_cleanup(context, api_key, reservation)
try:
result = await context.service.compact_responses(
payload,
Expand Down Expand Up @@ -5399,7 +5441,7 @@ async def _transcribe_request(
request_model=_TRANSCRIPTION_MODEL,
request_service_tier=None,
)
rate_limit_headers = await _rate_limit_headers_for_request(context, api_key)
rate_limit_headers = await _rate_limit_headers_with_reservation_cleanup(context, api_key, reservation)
try:
result = await context.service.transcribe(
audio_bytes=multipart.audio_bytes,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-07-31
Loading
Loading