From 7545d958d10943d44d946f9d4d51c8452076b88b Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 16:11:31 +0400 Subject: [PATCH 01/23] fix(proxy): retry silent bridge response.create upstreams --- .../_service/http_bridge/upstream_events.py | 34 +++--- .../integration/test_http_responses_bridge.py | 102 ++++++++++++++++++ 2 files changed, 122 insertions(+), 14 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index a70367af9b..72dbe46896 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -624,16 +624,6 @@ async def _relay_http_bridge_upstream_messages( if not expired_owner: continue pending_count = len(session.pending_requests) - for request_state in session.pending_requests: - if request_state.failure_phase_override is None: - request_state.failure_phase_override = "upstream" - if request_state.failure_detail_override is None: - request_state.failure_detail_override = ( - _HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL - ) - # Claim the session before cancelling receive so a - # gate waiter cannot reopen this ambiguous socket. - session.closed = True if receive_task is not None: receive_cancelled = await _cancel_http_bridge_reader_child( receive_task, @@ -641,10 +631,6 @@ async def _relay_http_bridge_upstream_messages( ) if receive_cancelled: receive_task = None - _record_http_bridge_stuck_retire( - reason=_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL, - session=session, - ) _log_http_bridge_event( "missing_response_created_timeout", session.key, @@ -657,6 +643,26 @@ async def _relay_http_bridge_upstream_messages( _extract_model_class(session.request_model) if session.request_model else None ), ) + retried = await self._retry_http_bridge_precreated_request(session) + if retried: + continue + # Claim the session only after the safe pre-created + # retry path refuses it. Reconnecting first prevents + # a silent upstream websocket from stranding clients + # until the request budget expires. + session.closed = True + _record_http_bridge_stuck_retire( + reason=_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL, + session=session, + ) + async with session.pending_lock: + for request_state in session.pending_requests: + if request_state.failure_phase_override is None: + request_state.failure_phase_override = "upstream" + if request_state.failure_detail_override is None: + request_state.failure_detail_override = ( + _HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL + ) await self._fail_http_bridge_reader_and_maybe_retire( session, error_code="upstream_request_timeout", diff --git a/tests/integration/test_http_responses_bridge.py b/tests/integration/test_http_responses_bridge.py index 989783c5af..ee727569ac 100644 --- a/tests/integration/test_http_responses_bridge.py +++ b/tests/integration/test_http_responses_bridge.py @@ -8178,6 +8178,108 @@ async def fake_connect_responses_websocket( assert connect_count == 2 +@pytest.mark.asyncio +async def test_v1_responses_http_bridge_retries_when_upstream_never_acknowledges_response_create( + async_client, + monkeypatch, +): + _install_bridge_settings_with_limits( + monkeypatch, + enabled=True, + stuck_gate_retire_after_seconds=0.01, + ) + account_id = await _import_account( + async_client, + "acc_http_bridge_missing_created_retry", + "http-bridge-missing-created-retry@example.com", + ) + account = await _get_account(account_id) + silent_upstream = _SilentUpstreamWebSocket() + recovered_upstream = _FakeBridgeUpstreamWebSocket() + upstreams = [silent_upstream, recovered_upstream] + connect_count = 0 + + async def fake_select_account_with_budget( + self, + deadline, + *, + request_id, + kind, + request_stage="first_turn", + sticky_key, + sticky_kind, + reallocate_sticky, + sticky_max_age_seconds, + prefer_earlier_reset_accounts, + routing_strategy, + model, + exclude_account_ids=None, + additional_limit_name=None, + api_key=None, + preferred_account_id=None, + ): + del preferred_account_id + del ( + self, + deadline, + request_id, + kind, + request_stage, + sticky_key, + sticky_kind, + reallocate_sticky, + sticky_max_age_seconds, + prefer_earlier_reset_accounts, + routing_strategy, + model, + exclude_account_ids, + additional_limit_name, + api_key, + ) + return AccountSelection(account=account, error_message=None, error_code=None) + + async def fake_ensure_fresh_with_budget(self, target, *, force=False, timeout_seconds): + del self, force, timeout_seconds + return target + + async def fake_connect_responses_websocket( + headers, + access_token, + account_id_header, + *, + base_url=None, + session=None, + ): + del headers, access_token, account_id_header, base_url, session + nonlocal connect_count + upstream = upstreams[connect_count] + connect_count += 1 + return upstream + + monkeypatch.setattr(proxy_module.ProxyService, "_select_account_with_budget", fake_select_account_with_budget) + monkeypatch.setattr(proxy_module.ProxyService, "_ensure_fresh_with_budget", fake_ensure_fresh_with_budget) + monkeypatch.setattr(proxy_module, "connect_responses_websocket", fake_connect_responses_websocket) + + response = await asyncio.wait_for( + async_client.post( + "/v1/responses", + json={ + "model": "gpt-5.1", + "instructions": "Return exactly OK.", + "input": "retry missing response.created", + "prompt_cache_key": "missing-created-retry-key", + }, + ), + timeout=_TEST_SYNC_TIMEOUT_SECONDS, + ) + + assert response.status_code == 200 + assert connect_count == 2 + assert silent_upstream.closed is True + assert len(silent_upstream.sent_text) == 1 + assert len(recovered_upstream.sent_text) == 1 + + @pytest.mark.asyncio async def test_backend_responses_http_bridge_retries_precreated_server_overload(async_client, monkeypatch): _install_bridge_settings(monkeypatch, enabled=True) From 09a73da551e8cc063fc13618818b6adb906ad92c Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 16:33:42 +0400 Subject: [PATCH 02/23] fix(proxy): retry bridge missing-created at deadline --- .../proxy/_service/http_bridge/request_submit.py | 12 ++++++++++++ .../proxy/_service/http_bridge/upstream_events.py | 5 ++++- tests/integration/test_http_responses_bridge.py | 2 +- tests/unit/test_proxy_http_bridge.py | 2 +- 4 files changed, 18 insertions(+), 3 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 0319e918ea..43f750e233 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -1530,12 +1530,14 @@ async def _retry_http_bridge_precreated_request( session: "_HTTPBridgeSession", *, request_state: _WebSocketRequestState | None = None, + allow_expired_deadline: bool = False, ) -> bool: account_neutral_recovery = is_http_bridge_account_neutral_replay( kind=session.key.affinity_kind, key=session.key.affinity_key, ) hard_owner_bound = _http_bridge_key_strength(session.key) == "hard" + now = _service_time().monotonic() async with session.pending_lock: if request_state is not None: if ( @@ -1544,6 +1546,11 @@ async def _retry_http_bridge_precreated_request( or request_state.draining_until_terminal or not _http_bridge_request_counts_against_queue(request_state) or not _websocket_request_can_replay_before_visible_output(request_state) + or ( + not allow_expired_deadline + and request_state.bridge_request_deadline is not None + and request_state.bridge_request_deadline <= now + ) ): return False else: @@ -1552,6 +1559,11 @@ async def _retry_http_bridge_precreated_request( for request_state in session.pending_requests if not request_state.draining_until_terminal and _websocket_request_can_replay_before_visible_output(request_state) + and ( + allow_expired_deadline + or request_state.bridge_request_deadline is None + or request_state.bridge_request_deadline > now + ) ] if len(retryable_requests) != 1: return False diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index 72dbe46896..a12deff177 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -643,7 +643,10 @@ async def _relay_http_bridge_upstream_messages( _extract_model_class(session.request_model) if session.request_model else None ), ) - retried = await self._retry_http_bridge_precreated_request(session) + retried = await self._retry_http_bridge_precreated_request( + session, + allow_expired_deadline=True, + ) if retried: continue # Claim the session only after the safe pre-created diff --git a/tests/integration/test_http_responses_bridge.py b/tests/integration/test_http_responses_bridge.py index ee727569ac..02d7713502 100644 --- a/tests/integration/test_http_responses_bridge.py +++ b/tests/integration/test_http_responses_bridge.py @@ -8186,8 +8186,8 @@ async def test_v1_responses_http_bridge_retries_when_upstream_never_acknowledges _install_bridge_settings_with_limits( monkeypatch, enabled=True, - stuck_gate_retire_after_seconds=0.01, ) + proxy_module.get_settings().http_responses_session_bridge_stuck_gate_retire_after_seconds = 0.01 account_id = await _import_account( async_client, "acc_http_bridge_missing_created_retry", diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 00dd232e8e..d74f1cbbaa 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -18444,7 +18444,7 @@ async def close(self) -> None: assert owner.response_event_count == 0 if leading_telemetry: assert owner.latency_first_upstream_event_ms is not None - retry_precreated.assert_not_awaited() + retry_precreated.assert_awaited_once_with(session, allow_expired_deadline=True) assert write_request_log.await_count == 2 assert {call.kwargs["error_code"] for call in write_request_log.await_args_list} == {"upstream_request_timeout"} fail_reader.assert_awaited_once() From f9ab3636dbab203b2b88b890d2a0d0eda7f94e5c Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 16:45:26 +0400 Subject: [PATCH 03/23] fix(proxy): retry same-anchor missing-created turns --- .../_service/http_bridge/request_submit.py | 54 ++++++++++++++++++- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 43f750e233..b54d86a54a 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -127,6 +127,7 @@ from app.modules.proxy._service.support import ( _HARD_HTTP_BRIDGE_AFFINITY_KINDS, # noqa: F401 _WEBSOCKET_FULL_REPLAY_WAIT_POLL_SECONDS, # noqa: F401 + _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS, _clear_websocket_request_error_overrides, _copy_websocket_route_metadata_from_session, _event_type_from_payload, @@ -199,6 +200,37 @@ ) +def _http_bridge_can_replay_same_anchor_before_created(request_state: _WebSocketRequestState) -> bool: + if not request_state.request_text: + return False + if request_state.replay_count >= _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS: + return False + return ( + request_state.previous_response_id is not None + and request_state.response_id is None + and request_state.awaiting_response_created + and request_state.response_event_count == 0 + and request_state.last_downstream_sequence_number is None + and not request_state.downstream_visible + and not request_state.upstream_model_output_seen + ) + + +async def _await_task_deferring_cancellation( + task: asyncio.Task[T], +) -> tuple[T, asyncio.CancelledError | None]: + """Finish critical cleanup while preserving the caller's cancellation.""" + + cancellation: asyncio.CancelledError | None = None + while True: + try: + return await asyncio.shield(task), cancellation + except asyncio.CancelledError as exc: + if task.cancelled(): + raise + cancellation = cancellation or exc + + async def _rollback_http_bridge_recovery_turn_state_registration( service: Any, receipt: DurableBridgeAliasRegistrationReceipt, @@ -1545,7 +1577,14 @@ async def _retry_http_bridge_precreated_request( or any(pending_request is not request_state for pending_request in session.pending_requests) or request_state.draining_until_terminal or not _http_bridge_request_counts_against_queue(request_state) - or not _websocket_request_can_replay_before_visible_output(request_state) + or not ( + _websocket_request_can_replay_before_visible_output(request_state) + or ( + allow_expired_deadline + and hard_owner_bound + and _http_bridge_can_replay_same_anchor_before_created(request_state) + ) + ) or ( not allow_expired_deadline and request_state.bridge_request_deadline is not None @@ -1558,7 +1597,14 @@ async def _retry_http_bridge_precreated_request( request_state for request_state in session.pending_requests if not request_state.draining_until_terminal - and _websocket_request_can_replay_before_visible_output(request_state) + and ( + _websocket_request_can_replay_before_visible_output(request_state) + or ( + allow_expired_deadline + and hard_owner_bound + and _http_bridge_can_replay_same_anchor_before_created(request_state) + ) + ) and ( allow_expired_deadline or request_state.bridge_request_deadline is None @@ -1572,6 +1618,10 @@ async def _retry_http_bridge_precreated_request( request_state.proxy_injected_previous_response_id and request_state.fresh_upstream_request_is_retry_safe and request_state.fresh_upstream_request_text + ) and not ( + allow_expired_deadline + and hard_owner_bound + and _http_bridge_can_replay_same_anchor_before_created(request_state) ): # Once a continuation is pending upstream, reconnecting without # replay cannot complete the current request, while replaying it From 98ef1aa394e6587ddc1901c799797a19a76324a5 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 19:05:34 +0400 Subject: [PATCH 04/23] fix(proxy): keep recovering silent response.create upstreams --- app/modules/proxy/_service/support.py | 1 + tests/integration/test_http_responses_bridge.py | 10 +++++----- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 0d842e78f9..c3133aa3ce 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -67,6 +67,7 @@ _TTFT_OUTPUT_ITEM_TYPES = _PENDING_TOOL_CALL_ITEM_TYPES - {"function_call"} _WEBSOCKET_FULL_REPLAY_WAIT_MIN_ITEMS = 20 _WEBSOCKET_FULL_REPLAY_WAIT_POLL_SECONDS = 0.05 +_WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS = 20 _HARD_HTTP_BRIDGE_AFFINITY_KINDS = frozenset( { "turn_state_header", diff --git a/tests/integration/test_http_responses_bridge.py b/tests/integration/test_http_responses_bridge.py index 02d7713502..b12770af52 100644 --- a/tests/integration/test_http_responses_bridge.py +++ b/tests/integration/test_http_responses_bridge.py @@ -8194,9 +8194,9 @@ async def test_v1_responses_http_bridge_retries_when_upstream_never_acknowledges "http-bridge-missing-created-retry@example.com", ) account = await _get_account(account_id) - silent_upstream = _SilentUpstreamWebSocket() + silent_upstreams = [_SilentUpstreamWebSocket() for _ in range(5)] recovered_upstream = _FakeBridgeUpstreamWebSocket() - upstreams = [silent_upstream, recovered_upstream] + upstreams = [*silent_upstreams, recovered_upstream] connect_count = 0 async def fake_select_account_with_budget( @@ -8274,9 +8274,9 @@ async def fake_connect_responses_websocket( ) assert response.status_code == 200 - assert connect_count == 2 - assert silent_upstream.closed is True - assert len(silent_upstream.sent_text) == 1 + assert connect_count == len(silent_upstreams) + 1 + assert all(upstream.closed for upstream in silent_upstreams) + assert [len(upstream.sent_text) for upstream in silent_upstreams] == [1] * len(silent_upstreams) assert len(recovered_upstream.sent_text) == 1 From cb1c3ce35094aa8c47cfbad2bbc78b07e00c4cf4 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 13:42:13 +0400 Subject: [PATCH 05/23] fix(proxy): retry silent bridge creates sooner --- app/modules/proxy/_service/http_bridge/helpers.py | 2 +- tests/unit/test_proxy_http_bridge.py | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/helpers.py b/app/modules/proxy/_service/http_bridge/helpers.py index ca98a7fa9f..0b28e0fd29 100644 --- a/app/modules/proxy/_service/http_bridge/helpers.py +++ b/app/modules/proxy/_service/http_bridge/helpers.py @@ -188,7 +188,7 @@ logger = logging.getLogger("app.modules.proxy.service") _HTTP_BRIDGE_BACKGROUND_CLOSE_TIMEOUT_SECONDS = 5.0 -_HTTP_BRIDGE_EVENTLESS_RESPONSE_CREATED_MAX_SECONDS = 240.0 +_HTTP_BRIDGE_EVENTLESS_RESPONSE_CREATED_MAX_SECONDS = 15.0 _HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL = "missing_response_created_timeout" T = TypeVar("T") diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index d74f1cbbaa..72e60a9bc7 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -136,14 +136,14 @@ def test_http_bridge_eventless_precreated_deadline_uses_current_send_and_client_ request_state, stuck_gate_retire_after_seconds=300.0, ) - == 340.0 + == 115.0 ) assert ( http_bridge_helpers_module._http_bridge_eventless_precreated_deadline( request_state, - stuck_gate_retire_after_seconds=30.0, + stuck_gate_retire_after_seconds=10.0, ) - == 130.0 + == 110.0 ) request_state.latency_first_upstream_event_ms = 25 @@ -152,7 +152,7 @@ def test_http_bridge_eventless_precreated_deadline_uses_current_send_and_client_ request_state, stuck_gate_retire_after_seconds=300.0, ) - == 340.0 + == 115.0 ) From 1bc93d6dccd19d0abdc70505f75854c6b09c22cf Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 14:13:09 +0400 Subject: [PATCH 06/23] fix(proxy): bound missing-created bridge retries --- .../_service/http_bridge/request_submit.py | 22 +++++++-- app/modules/proxy/_service/support.py | 1 + tests/unit/test_proxy_http_bridge.py | 49 +++++++++++++++++++ 3 files changed, 69 insertions(+), 3 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index b54d86a54a..7645ad4d28 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -193,6 +193,7 @@ _REQUEST_TRANSPORT_HTTP = "http" _WEBSOCKET_AUTH_INVALIDATED_FAILURE_CODE = "account_auth_invalidated" _NO_SECURITY_WORK_AUTHORIZED_ACCOUNTS_CODE = "no_security_work_authorized_accounts" +_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_MAX_RETRIES = 1 _SECURITY_WORK_NO_AUTHORIZED_ACCOUNTS_MESSAGE = ( "Upstream flagged this request as possible cybersecurity work, but no account is marked as authorized for " "security work. codex-lb is continuing with normal account selection; the upstream request may still fail until " @@ -216,6 +217,13 @@ def _http_bridge_can_replay_same_anchor_before_created(request_state: _WebSocket ) +def _http_bridge_can_retry_missing_response_created(request_state: _WebSocketRequestState) -> bool: + return ( + request_state.missing_response_created_retry_count < _HTTP_BRIDGE_MISSING_RESPONSE_CREATED_MAX_RETRIES + and _http_bridge_can_replay_same_anchor_before_created(request_state) + ) + + async def _await_task_deferring_cancellation( task: asyncio.Task[T], ) -> tuple[T, asyncio.CancelledError | None]: @@ -1570,6 +1578,7 @@ async def _retry_http_bridge_precreated_request( ) hard_owner_bound = _http_bridge_key_strength(session.key) == "hard" now = _service_time().monotonic() + missing_created_retry = False async with session.pending_lock: if request_state is not None: if ( @@ -1582,7 +1591,7 @@ async def _retry_http_bridge_precreated_request( or ( allow_expired_deadline and hard_owner_bound - and _http_bridge_can_replay_same_anchor_before_created(request_state) + and _http_bridge_can_retry_missing_response_created(request_state) ) ) or ( @@ -1602,7 +1611,7 @@ async def _retry_http_bridge_precreated_request( or ( allow_expired_deadline and hard_owner_bound - and _http_bridge_can_replay_same_anchor_before_created(request_state) + and _http_bridge_can_retry_missing_response_created(request_state) ) ) and ( @@ -1621,7 +1630,7 @@ async def _retry_http_bridge_precreated_request( ) and not ( allow_expired_deadline and hard_owner_bound - and _http_bridge_can_replay_same_anchor_before_created(request_state) + and _http_bridge_can_retry_missing_response_created(request_state) ): # Once a continuation is pending upstream, reconnecting without # replay cannot complete the current request, while replaying it @@ -1629,6 +1638,11 @@ async def _retry_http_bridge_precreated_request( # injected retry-safe anchors are equivalent to the client's own # full resend once the anchor is stripped. return False + missing_created_retry = ( + allow_expired_deadline + and hard_owner_bound + and _http_bridge_can_retry_missing_response_created(request_state) + ) close_classification = _classify_upstream_close( session.last_upstream_close_code, response_events_seen=request_state.response_event_count, @@ -1678,6 +1692,8 @@ async def _retry_http_bridge_precreated_request( elif not request_state.file_required_preferred_account and not hard_owner_bound: request_state.preferred_account_id = None request_state.excluded_account_ids.add(session.account.id) + if missing_created_retry: + request_state.missing_response_created_retry_count += 1 if session.account.id in request_state.excluded_account_ids: session.upstream_turn_state = None session.downstream_turn_state = None diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index c3133aa3ce..43b886ce7c 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -787,6 +787,7 @@ class _WebSocketRequestState: request_usage_budget: ApiKeyRequestUsageBudget | None = None request_text: str | None = None replay_count: int = 0 + missing_response_created_retry_count: int = 0 auth_replay_count: int = 0 auth_replay_counts_by_account: dict[str, int] = field(default_factory=dict) force_refresh_account_id: str | None = None diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 72e60a9bc7..6c03febe07 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -268,6 +268,55 @@ async def send_text(_text: str) -> None: ) +@pytest.mark.asyncio +async def test_http_bridge_missing_response_created_retries_once_before_terminal( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + request_text = ( + '{"type":"response.create","model":"gpt-5.6-sol","previous_response_id":"resp-owner-anchor","input":"hello"}' + ) + request_state = _make_eventless_http_bridge_owner(request_id="req-missing-created-once") + request_state.previous_response_id = "resp-owner-anchor" + request_state.request_text = request_text + request_state.bridge_request_deadline = time.monotonic() - 1.0 + session = _make_bridge_session( + key_value="missing-created-once", + pending_requests=deque([request_state]), + queued_request_count=1, + ) + send_text = AsyncMock() + session.upstream = cast( + UpstreamWebSocket, + SimpleNamespace(send_text=send_text, close=AsyncMock()), + ) + reconnect = AsyncMock() + monkeypatch.setattr(service, "_reconnect_http_bridge_session", reconnect) + monkeypatch.setattr(service, "_acquire_account_response_create_lease_or_overload", AsyncMock(return_value=None)) + + first_retry = await service._retry_http_bridge_precreated_request( + session, + allow_expired_deadline=True, + ) + second_retry = await service._retry_http_bridge_precreated_request( + session, + allow_expired_deadline=True, + ) + + assert first_retry is True + assert second_retry is False + reconnect.assert_awaited_once_with( + session, + request_state=request_state, + require_same_account=True, + ) + send_text.assert_awaited_once() + assert json.loads(send_text.await_args.args[0])["previous_response_id"] == "resp-owner-anchor" + assert request_state.missing_response_created_retry_count == 1 + assert request_state.replay_count == 1 + assert request_state.awaiting_response_created is True + + def _make_account_neutral_replay_session_key( nonce: str, api_key_id: str | None = None, From ad28520c0ac64558d62ed0d452fed769f4648155 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 14:40:17 +0400 Subject: [PATCH 07/23] fix(proxy): penalize silent bridge create owners --- app/modules/proxy/_service/http_bridge/upstream_events.py | 2 +- tests/unit/test_proxy_http_bridge.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index a12deff177..f09b8654b7 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -670,7 +670,7 @@ async def _relay_http_bridge_upstream_messages( session, error_code="upstream_request_timeout", error_message=receive_timeout.error_message, - penalize_account=False, + penalize_account=True, retire_detail=_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL, force_retire=True, ) diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 6c03febe07..4aa0912cd2 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -18497,7 +18497,7 @@ async def close(self) -> None: assert write_request_log.await_count == 2 assert {call.kwargs["error_code"] for call in write_request_log.await_args_list} == {"upstream_request_timeout"} fail_reader.assert_awaited_once() - assert fail_reader.await_args.kwargs["penalize_account"] is False + assert fail_reader.await_args.kwargs["penalize_account"] is True assert fail_reader.await_args.kwargs["force_retire"] is True record_stuck_retire.assert_called_once_with( reason="missing_response_created_timeout", From 0115eab2fc69932ed4713e83cc0855151731f863 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 14:53:05 +0400 Subject: [PATCH 08/23] fix(proxy): rebind silent bridge create owners --- .../_service/http_bridge/request_submit.py | 36 ++++++++++- tests/unit/test_proxy_http_bridge.py | 62 +++++++++++++++++++ 2 files changed, 96 insertions(+), 2 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 7645ad4d28..09afae3841 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -1579,6 +1579,9 @@ async def _retry_http_bridge_precreated_request( hard_owner_bound = _http_bridge_key_strength(session.key) == "hard" now = _service_time().monotonic() missing_created_retry = False + rebind_missing_created_owner = False + owner_rebind_affinity: _AffinityPolicy | None = None + selection_rebind_affinity: _AffinityPolicy | None = None async with session.pending_lock: if request_state is not None: if ( @@ -1680,7 +1683,25 @@ async def _retry_http_bridge_precreated_request( request_text = _prepare_websocket_request_state_for_visible_output_replay(request_state) if request_text is None: return False - if not hard_owner_bound: + if missing_created_retry and hard_owner_bound: + request_state.preferred_account_id = None + request_state.excluded_account_ids.add(session.account.id) + request_state.affinity_policy = replace( + request_state.affinity_policy, + key=None, + kind=None, + reallocate_sticky=True, + ) + rebind_missing_created_owner = True + owner_rebind_affinity = session.affinity + selection_rebind_affinity = replace( + session.affinity, + key=None, + kind=None, + reallocate_sticky=True, + codex_session_source=None, + ) + elif not hard_owner_bound: request_state.excluded_account_ids.add(session.account.id) else: require_preferred_reconnect = account_neutral_recovery @@ -1710,7 +1731,18 @@ async def _retry_http_bridge_precreated_request( model_class=_extract_model_class(session.request_model) if session.request_model else None, ) try: - if hard_owner_bound: + if rebind_missing_created_owner: + await _call_with_supported_optional_kwargs( + self._reconnect_http_bridge_session, + session, + optional_kwargs={ + "owner_rebind_affinity": owner_rebind_affinity, + "selection_affinity": selection_rebind_affinity, + }, + request_state=request_state, + require_same_account=account_neutral_recovery, + ) + elif hard_owner_bound: await self._reconnect_http_bridge_session( session, request_state=request_state, diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 4aa0912cd2..96df27c7c8 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -317,6 +317,68 @@ async def test_http_bridge_missing_response_created_retries_once_before_terminal assert request_state.awaiting_response_created is True +@pytest.mark.asyncio +async def test_http_bridge_missing_response_created_rebinds_hard_owner_when_full_resend_is_safe( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + anchored_text = ( + '{"type":"response.create","model":"gpt-5.6-sol","previous_response_id":"resp-owner-anchor","input":"trimmed"}' + ) + fresh_text = ( + '{"type":"response.create","model":"gpt-5.6-sol",' + '"input":[{"role":"user","content":[{"type":"input_text","text":"hello"}]}]}' + ) + request_state = _make_eventless_http_bridge_owner(request_id="req-missing-created-rebind") + request_state.previous_response_id = "resp-owner-anchor" + request_state.proxy_injected_previous_response_id = True + request_state.fresh_upstream_request_is_retry_safe = True + request_state.fresh_upstream_request_text = fresh_text + request_state.request_text = anchored_text + request_state.bridge_request_deadline = time.monotonic() - 1.0 + session = _make_bridge_session( + key=proxy_service._HTTPBridgeSessionKey("session_header", "missing-created-rebind", None), + key_value="missing-created-rebind", + pending_requests=deque([request_state]), + queued_request_count=1, + ) + original_affinity = session.affinity + send_text = AsyncMock() + session.upstream = cast( + UpstreamWebSocket, + SimpleNamespace(send_text=send_text, close=AsyncMock()), + ) + reconnect = AsyncMock() + monkeypatch.setattr(service, "_reconnect_http_bridge_session", reconnect) + monkeypatch.setattr(service, "_acquire_account_response_create_lease_or_overload", AsyncMock(return_value=None)) + + retried = await service._retry_http_bridge_precreated_request( + session, + allow_expired_deadline=True, + ) + + assert retried is True + reconnect.assert_awaited_once() + reconnect_kwargs = reconnect.await_args.kwargs + assert reconnect_kwargs["request_state"] is request_state + assert reconnect_kwargs["require_same_account"] is False + assert reconnect_kwargs["owner_rebind_affinity"] is original_affinity + selection_affinity = reconnect_kwargs["selection_affinity"] + assert selection_affinity.key is None + assert selection_affinity.kind is None + assert selection_affinity.reallocate_sticky is True + send_text.assert_awaited_once() + assert json.loads(send_text.await_args.args[0]).get("previous_response_id") is None + assert request_state.request_text == fresh_text + assert request_state.previous_response_id is None + assert request_state.preferred_account_id is None + assert request_state.excluded_account_ids == {"acc-bridge"} + assert request_state.affinity_policy.reallocate_sticky is True + assert request_state.missing_response_created_retry_count == 1 + assert request_state.replay_count == 1 + assert request_state.awaiting_response_created is True + + def _make_account_neutral_replay_session_key( nonce: str, api_key_id: str | None = None, From a6347cc99699949542f27f652471b94476d0bed3 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 15:05:29 +0400 Subject: [PATCH 09/23] fix(proxy): clear silent bridge anchors on terminal retire --- .../proxy/_service/http_bridge/helpers.py | 3 +- .../proxy/_service/http_bridge/mixin.py | 10 +++- .../_service/http_bridge/request_submit.py | 6 ++- .../proxy/durable_bridge_coordinator.py | 2 + .../proxy/durable_bridge_repository.py | 11 +++++ tests/unit/test_durable_bridge_sessions.py | 48 +++++++++++++++++++ tests/unit/test_proxy_http_bridge.py | 35 ++++++++++++++ 7 files changed, 112 insertions(+), 3 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/helpers.py b/app/modules/proxy/_service/http_bridge/helpers.py index 0b28e0fd29..b2e7494944 100644 --- a/app/modules/proxy/_service/http_bridge/helpers.py +++ b/app/modules/proxy/_service/http_bridge/helpers.py @@ -680,11 +680,12 @@ async def _close_http_bridge_session_bounded( session: "_HTTPBridgeSession", *, reason: str, + clear_continuity: bool = False, ) -> None: if session.upstream_reader is asyncio.current_task(): session.upstream_reader = None close_task = asyncio.create_task( - service._close_http_bridge_session(session), + service._close_http_bridge_session(session, clear_continuity=clear_continuity), name=f"http-bridge-close-{_hash_identifier(session.key.affinity_key)}", ) diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index d8dfeb0c51..dde849131e 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -242,8 +242,14 @@ async def _close_http_bridge_session_bounded( session: "_HTTPBridgeSession", *, reason: str, + clear_continuity: bool = False, ) -> None: - await _close_http_bridge_session_bounded(self, session, reason=reason) + await _close_http_bridge_session_bounded( + self, + session, + reason=reason, + clear_continuity=clear_continuity, + ) def _schedule_http_bridge_session_closes( self, @@ -1641,6 +1647,7 @@ async def _close_http_bridge_session( session: "_HTTPBridgeSession", *, turn_state_lock_held: bool = False, + clear_continuity: bool = False, ) -> None: session.closed = True if turn_state_lock_held: @@ -1663,6 +1670,7 @@ async def _close_http_bridge_session( instance_id=_service_get_settings().http_responses_session_bridge_instance_id, owner_epoch=session.durable_owner_epoch, draining=shutdown_state.is_bridge_drain_active(), + clear_continuity=clear_continuity, ) except Exception: logger.warning("Failed to release durable HTTP bridge session", exc_info=True) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 09afae3841..e711344a7b 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -1476,7 +1476,11 @@ async def _retire_stale_pending_http_bridge_session( if should_close: session.upstream_close_attempted = True if should_close: - await self._close_http_bridge_session_bounded(session, reason="retire_stale_pending") + await self._close_http_bridge_session_bounded( + session, + reason="retire_stale_pending", + clear_continuity=detail == "missing_response_created_timeout", + ) _log_http_bridge_event( "retire_stale_pending", session.key, diff --git a/app/modules/proxy/durable_bridge_coordinator.py b/app/modules/proxy/durable_bridge_coordinator.py index 75e397a7d2..263520ca3b 100644 --- a/app/modules/proxy/durable_bridge_coordinator.py +++ b/app/modules/proxy/durable_bridge_coordinator.py @@ -268,6 +268,7 @@ async def release_live_session( instance_id: str, owner_epoch: int, draining: bool, + clear_continuity: bool = False, ) -> DurableBridgeLookup | None: async with self._session() as session: snapshot = await DurableBridgeRepository(session).release_session( @@ -275,6 +276,7 @@ async def release_live_session( instance_id=instance_id, owner_epoch=owner_epoch, draining=draining, + clear_continuity=clear_continuity, ) if snapshot is None: return None diff --git a/app/modules/proxy/durable_bridge_repository.py b/app/modules/proxy/durable_bridge_repository.py index 7b82e7f8a3..be501270ba 100644 --- a/app/modules/proxy/durable_bridge_repository.py +++ b/app/modules/proxy/durable_bridge_repository.py @@ -385,6 +385,7 @@ async def release_session( instance_id: str, owner_epoch: int, draining: bool, + clear_continuity: bool = False, ) -> DurableBridgeSessionSnapshot | None: """Release the lease with a single fenced UPDATE. @@ -399,6 +400,16 @@ async def release_session( "state": HttpBridgeSessionState.DRAINING if draining else HttpBridgeSessionState.CLOSED, "closed_at": None if draining else now, } + if clear_continuity: + values.update( + { + "latest_turn_state": None, + "latest_response_id": None, + "latest_input_item_count": None, + "latest_input_full_fingerprint": None, + "latest_pending_tool_calls_json": None, + } + ) return await self._execute_fenced_session_update( session_id=session_id, instance_id=instance_id, diff --git a/tests/unit/test_durable_bridge_sessions.py b/tests/unit/test_durable_bridge_sessions.py index 7d52d9376d..6338791e84 100644 --- a/tests/unit/test_durable_bridge_sessions.py +++ b/tests/unit/test_durable_bridge_sessions.py @@ -1760,6 +1760,54 @@ async def test_durable_bridge_same_account_closed_takeover_preserves_restart_anc assert reclaimed.latest_response_id == "resp_old" +@pytest.mark.asyncio +async def test_durable_bridge_terminal_release_can_clear_restart_anchor( + coordinator: DurableBridgeSessionCoordinator, +) -> None: + claimed = await coordinator.claim_live_session( + session_key_kind="session_header", + session_key_value="sid-terminal-clear", + api_key_id=None, + instance_id="instance-a", + lease_ttl_seconds=60.0, + account_id="acc-1", + model="gpt-5.4", + service_tier=None, + latest_turn_state="http_turn_old", + latest_response_id="resp_old", + allow_takeover=True, + ) + released = await coordinator.release_live_session( + session_id=claimed.session_id, + instance_id="instance-a", + owner_epoch=claimed.owner_epoch, + draining=False, + clear_continuity=True, + ) + + assert released is not None + assert released.latest_turn_state is None + assert released.latest_response_id is None + + reclaimed = await coordinator.claim_live_session( + session_key_kind="session_header", + session_key_value="sid-terminal-clear", + api_key_id=None, + instance_id="instance-b", + lease_ttl_seconds=60.0, + account_id="acc-1", + model="gpt-5.4", + service_tier=None, + latest_turn_state=None, + latest_response_id=None, + allow_takeover=True, + ) + + assert reclaimed.owner_instance_id == "instance-b" + assert reclaimed.latest_turn_state is None + assert reclaimed.latest_response_id is None + + @pytest.mark.asyncio async def test_durable_bridge_takeover_preserves_existing_anchor_when_replacement_has_none( coordinator: DurableBridgeSessionCoordinator, diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 96df27c7c8..e8dfa40fa3 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -15417,6 +15417,7 @@ async def test_http_bridge_retire_after_drain_waits_for_queued_submission( instance_id="instance-retire-drain", owner_epoch=3, draining=False, + clear_continuity=False, ) release_account_lease.assert_awaited_once_with(lease) assert session.account_lease is None @@ -15482,6 +15483,7 @@ async def test_http_bridge_retire_after_drain_does_not_cancel_current_upstream_r instance_id="instance-reader-retire", owner_epoch=7, draining=False, + clear_continuity=False, ) release_account_lease.assert_awaited_once_with(lease) assert session.account_lease is None @@ -18993,12 +18995,43 @@ async def test_retire_stale_pending_http_bridge_session_unregisters_aliases_and_ instance_id="instance-cleanup", owner_epoch=7, draining=False, + clear_continuity=False, ) release_account_lease.assert_awaited_once_with(lease) assert session.account_lease is None close.assert_awaited_once() +@pytest.mark.asyncio +async def test_retire_missing_created_http_bridge_session_clears_durable_continuity( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + release_live_session = AsyncMock() + service._durable_bridge = cast( + Any, + SimpleNamespace(release_live_session=release_live_session), + ) + session = _make_bridge_session(key_value="bridge-missing-created-cleanup") + session.durable_session_id = "durable-missing-created" + session.durable_owner_epoch = 9 + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-missing-created"), + ) + + await service._retire_stale_pending_http_bridge_session(session, detail="missing_response_created_timeout") + + release_live_session.assert_awaited_once_with( + session_id="durable-missing-created", + instance_id="instance-missing-created", + owner_epoch=9, + draining=False, + clear_continuity=True, + ) + + @pytest.mark.asyncio async def test_http_bridge_reader_failed_precreated_replay_retires_registered_session( monkeypatch: pytest.MonkeyPatch, @@ -19634,6 +19667,7 @@ async def fail_replay(target_session: proxy_service._HTTPBridgeSession) -> bool: instance_id="instance-log-fails", owner_epoch=3, draining=False, + clear_continuity=False, ) cast(Any, upstream).close.assert_awaited_once() @@ -19715,6 +19749,7 @@ async def test_http_bridge_reader_unexpected_processing_error_fails_pending_requ instance_id="instance-reader-crash", owner_epoch=9, draining=False, + clear_continuity=False, ) cast(Any, upstream).close.assert_awaited_once() write_request_log.assert_awaited_once() From 81c4f5310a975f579e749d4920718e95e6ab573b Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 15:25:28 +0400 Subject: [PATCH 10/23] fix(proxy): clear stale bridge anchors after close attempt --- .../_service/http_bridge/request_submit.py | 15 ++++++++- tests/unit/test_proxy_http_bridge.py | 33 +++++++++++++++++++ 2 files changed, 47 insertions(+), 1 deletion(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index e711344a7b..163532d843 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -10,6 +10,7 @@ import anyio +from app.core import shutdown as shutdown_state from app.core.clients.files import create_file as core_create_file # noqa: F401 from app.core.clients.files import finalize_file as core_finalize_file # noqa: F401 from app.core.clients.proxy import CodexControlResponse as CodexControlResponse @@ -1465,6 +1466,7 @@ async def _retire_stale_pending_http_bridge_session( *, detail: str, ) -> None: + clear_continuity = detail == "missing_response_created_timeout" session.closed = True async with self._http_bridge_lock: if self._http_bridge_sessions.get(session.key) is session: @@ -1479,8 +1481,19 @@ async def _retire_stale_pending_http_bridge_session( await self._close_http_bridge_session_bounded( session, reason="retire_stale_pending", - clear_continuity=detail == "missing_response_created_timeout", + clear_continuity=clear_continuity, ) + elif clear_continuity and session.durable_session_id is not None and session.durable_owner_epoch is not None: + try: + await self._durable_bridge.release_live_session( + session_id=session.durable_session_id, + instance_id=_service_get_settings().http_responses_session_bridge_instance_id, + owner_epoch=session.durable_owner_epoch, + draining=shutdown_state.is_bridge_drain_active(), + clear_continuity=True, + ) + except Exception: + logger.warning("Failed to clear durable HTTP bridge continuity during stale retire", exc_info=True) _log_http_bridge_event( "retire_stale_pending", session.key, diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index e8dfa40fa3..52e77b6078 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -19032,6 +19032,39 @@ async def test_retire_missing_created_http_bridge_session_clears_durable_continu ) +@pytest.mark.asyncio +async def test_retire_missing_created_http_bridge_session_clears_durable_continuity_after_close_attempt( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + release_live_session = AsyncMock() + service._durable_bridge = cast( + Any, + SimpleNamespace(release_live_session=release_live_session), + ) + session = _make_bridge_session(key_value="bridge-missing-created-after-close") + session.durable_session_id = "durable-missing-created-after-close" + session.durable_owner_epoch = 11 + session.upstream_close_attempted = True + close = cast(Any, session.upstream).close + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-missing-created-after-close"), + ) + + await service._retire_stale_pending_http_bridge_session(session, detail="missing_response_created_timeout") + + close.assert_not_awaited() + release_live_session.assert_awaited_once_with( + session_id="durable-missing-created-after-close", + instance_id="instance-missing-created-after-close", + owner_epoch=11, + draining=False, + clear_continuity=True, + ) + + @pytest.mark.asyncio async def test_http_bridge_reader_failed_precreated_replay_retires_registered_session( monkeypatch: pytest.MonkeyPatch, From 6e6045417291ccb0b4b3993eaff3b6870c748dc8 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 19:28:00 +0400 Subject: [PATCH 11/23] fix(proxy): limit created-only close replay budget --- app/modules/proxy/_service/support.py | 7 +++++++ tests/integration/test_proxy_websocket_responses.py | 2 +- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 43b886ce7c..197fc2a694 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -67,6 +67,7 @@ _TTFT_OUTPUT_ITEM_TYPES = _PENDING_TOOL_CALL_ITEM_TYPES - {"function_call"} _WEBSOCKET_FULL_REPLAY_WAIT_MIN_ITEMS = 20 _WEBSOCKET_FULL_REPLAY_WAIT_POLL_SECONDS = 0.05 +_WEBSOCKET_CREATED_ONLY_CLOSE_MAX_REPLAYS = 1 _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS = 20 _HARD_HTTP_BRIDGE_AFFINITY_KINDS = frozenset( { @@ -1169,6 +1170,12 @@ def _websocket_request_can_replay_before_visible_output(request_state: _WebSocke and request_state.response_event_count == 1 and not request_state.downstream_visible ) + if ( + request_state.response_id is not None + and not sequenced_created_only_prewarm + and request_state.replay_count >= _WEBSOCKET_CREATED_ONLY_CLOSE_MAX_REPLAYS + ): + return False if request_state.last_downstream_sequence_number is not None and not sequenced_created_only_prewarm: return False if request_state.downstream_visible: diff --git a/tests/integration/test_proxy_websocket_responses.py b/tests/integration/test_proxy_websocket_responses.py index 43c6b876a2..d840591381 100644 --- a/tests/integration/test_proxy_websocket_responses.py +++ b/tests/integration/test_proxy_websocket_responses.py @@ -9250,7 +9250,7 @@ async def fake_write_request_log(self, **kwargs): assert failed_event["response"]["error"]["code"] == "stream_incomplete" assert "close_code=1011" in failed_event["response"]["error"]["message"] assert len(log_calls) == 1 - assert log_calls[0]["request_id"] == "resp_ws_eof_retry" + assert log_calls[0]["request_id"] == "resp_ws_eof_retry_1" assert log_calls[0]["status"] == "error" assert log_calls[0]["error_code"] == "stream_incomplete" From 91b783ce2041f8b082574b1c62ada536ed35e03d Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 19:41:24 +0400 Subject: [PATCH 12/23] fix(proxy): retry previsible stream eof continuations --- app/modules/proxy/_service/streaming/mixin.py | 3 + tests/unit/test_proxy_utils.py | 61 +++++++++++++++++++ 2 files changed, 64 insertions(+) diff --git a/app/modules/proxy/_service/streaming/mixin.py b/app/modules/proxy/_service/streaming/mixin.py index 6189ae7cd5..9ccc258d8e 100644 --- a/app/modules/proxy/_service/streaming/mixin.py +++ b/app/modules/proxy/_service/streaming/mixin.py @@ -491,6 +491,7 @@ async def _stream_once( enforce_openai_sdk_contract: bool = True, ) -> AsyncIterator[str]: proxy = cast(_StreamingServiceProtocol, self) + settlement.reset() account_id_value = account.id access_token = proxy._encryptor.decrypt(account.access_token_encrypted) account_id = _header_account_id(account.chatgpt_account_id) @@ -577,6 +578,8 @@ async def _stream_once( settlement.record_success = False settlement.account_health_error = True settlement.error = {"message": error_message} + if allow_transient_retry: + raise _TransientStreamError(error_code, settlement.error) yield format_sse_event( response_failed_event( error_code, diff --git a/tests/unit/test_proxy_utils.py b/tests/unit/test_proxy_utils.py index 31910f68ca..69939bb96a 100644 --- a/tests/unit/test_proxy_utils.py +++ b/tests/unit/test_proxy_utils.py @@ -30073,6 +30073,67 @@ async def fake_stream(payload, headers, access_token, account_id, base_url=None, record_success.assert_not_awaited() +@pytest.mark.asyncio +async def test_stream_previsible_core_eof_with_previous_response_id_retries(monkeypatch): + settings = _make_proxy_settings() + request_logs = _RequestLogsRecorder() + service = proxy_service.ProxyService(_repo_factory(request_logs)) + account = _make_account("acc_previsible_core_eof") + request_logs.response_owner_by_id[("resp_parent", None, "sid-stream")] = account.id + handle_stream_error = AsyncMock() + record_success = AsyncMock() + stream_calls = 0 + + monkeypatch.setattr(proxy_service, "get_settings_cache", lambda: _SettingsCache(settings)) + monkeypatch.setattr(proxy_service, "get_settings", lambda: settings) + monkeypatch.setattr(proxy_service, "_MAX_TRANSIENT_SAME_ACCOUNT_RETRIES", 3) + monkeypatch.setattr(streaming_retry_module.ProcessNetworkRecovery, "wait", AsyncMock(return_value=None)) + monkeypatch.setattr(streaming_retry_module.asyncio, "sleep", AsyncMock()) + monkeypatch.setattr( + service._load_balancer, + "select_account", + AsyncMock(return_value=AccountSelection(account=account, error_message=None)), + ) + monkeypatch.setattr(service._load_balancer, "record_success", record_success) + monkeypatch.setattr(service, "_ensure_fresh", AsyncMock(return_value=account)) + monkeypatch.setattr(service, "_handle_stream_error", handle_stream_error) + monkeypatch.setattr(service, "_settle_stream_api_key_usage", AsyncMock(return_value=True)) + + async def fake_stream(payload, headers, access_token, account_id, base_url=None, raise_for_status=False, **kwargs): + nonlocal stream_calls + del payload, headers, access_token, account_id, base_url, raise_for_status, kwargs + stream_calls += 1 + if stream_calls == 1: + return + yield 'data: {"type":"response.completed","response":{"id":"resp_child_retry_ok"}}\n\n' + + monkeypatch.setattr(proxy_service, "core_stream_responses", fake_stream) + + payload = ResponsesRequest.model_validate( + { + "model": "gpt-5.1", + "instructions": "hi", + "input": [], + "stream": True, + "previous_response_id": "resp_parent", + } + ) + + chunks = [chunk async for chunk in service.stream_responses(payload, {"session_id": "sid-stream"})] + + completed = json.loads(chunks[-1].split("data: ", 1)[1]) + assert completed["type"] == "response.completed" + assert completed["response"]["id"] == "resp_child_retry_ok" + assert stream_calls == 2 + assert request_logs.lookup_calls == [("resp_parent", None, "sid-stream")] + assert await service.drain_persistence_tasks(timeout_seconds=1) + assert [call["status"] for call in request_logs.calls] == ["error", "success"] + assert request_logs.calls[0]["error_code"] == "stream_incomplete" + assert request_logs.calls[-1]["request_id"] == "resp_child_retry_ok" + handle_stream_error.assert_not_awaited() + record_success.assert_awaited_once_with(account) + + @pytest.mark.asyncio async def test_stream_missing_tool_output_proxy_error_is_masked_to_stream_incomplete(monkeypatch, caplog): settings = _make_proxy_settings() From 1387a5d97d739ee207a8be3bbd4fd84cd65f900f Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Thu, 30 Jul 2026 20:11:36 +0400 Subject: [PATCH 13/23] fix(proxy): keep fresh empty streams fail-closed --- app/modules/proxy/_service/streaming/mixin.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/app/modules/proxy/_service/streaming/mixin.py b/app/modules/proxy/_service/streaming/mixin.py index 9ccc258d8e..a8d228cab8 100644 --- a/app/modules/proxy/_service/streaming/mixin.py +++ b/app/modules/proxy/_service/streaming/mixin.py @@ -298,6 +298,7 @@ _RetryableStreamError, _StreamSettlement, _TerminalStreamError, + _TransientStreamError, _ttft_event_latency_ms, _WebSocketUpstreamControl, ) @@ -578,7 +579,7 @@ async def _stream_once( settlement.record_success = False settlement.account_health_error = True settlement.error = {"message": error_message} - if allow_transient_retry: + if allow_transient_retry and payload.previous_response_id is not None: raise _TransientStreamError(error_code, settlement.error) yield format_sse_event( response_failed_event( From 2b9378de9a90010462a8efe882bf8b820d77b5ae Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 17:06:09 +0400 Subject: [PATCH 14/23] fix(proxy): keep retrying unanchored silent bridge creates --- app/modules/proxy/_service/support.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 197fc2a694..0afbddd839 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -1161,7 +1161,16 @@ def _websocket_request_can_replay_before_visible_output(request_state: _WebSocke if not request_state.request_text: return False if request_state.replay_count >= 1: - return False + unanchored_precreated_pending = ( + request_state.previous_response_id is None + and request_state.response_id is None + and request_state.awaiting_response_created + ) + if ( + not unanchored_precreated_pending + or request_state.replay_count >= _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS + ): + return False sequenced_created_only_prewarm = ( request_state.generate_false_prewarm and request_state.last_downstream_sequence_number == 0 From a448b9cae7b79f634a5b3c35a75d7bec137d3bae Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 17:07:05 +0400 Subject: [PATCH 15/23] fix(proxy): use shared bridge cancellation helper --- .../_service/http_bridge/request_submit.py | 35 +++++++------------ app/modules/proxy/_service/support.py | 5 +-- 2 files changed, 13 insertions(+), 27 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 163532d843..7d64b8de0f 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -225,21 +225,6 @@ def _http_bridge_can_retry_missing_response_created(request_state: _WebSocketReq ) -async def _await_task_deferring_cancellation( - task: asyncio.Task[T], -) -> tuple[T, asyncio.CancelledError | None]: - """Finish critical cleanup while preserving the caller's cancellation.""" - - cancellation: asyncio.CancelledError | None = None - while True: - try: - return await asyncio.shield(task), cancellation - except asyncio.CancelledError as exc: - if task.cancelled(): - raise - cancellation = cancellation or exc - - async def _rollback_http_bridge_recovery_turn_state_registration( service: Any, receipt: DurableBridgeAliasRegistrationReceipt, @@ -1643,14 +1628,18 @@ async def _retry_http_bridge_precreated_request( if len(retryable_requests) != 1: return False request_state = retryable_requests[0] - if request_state.previous_response_id is not None and not ( - request_state.proxy_injected_previous_response_id - and request_state.fresh_upstream_request_is_retry_safe - and request_state.fresh_upstream_request_text - ) and not ( - allow_expired_deadline - and hard_owner_bound - and _http_bridge_can_retry_missing_response_created(request_state) + if ( + request_state.previous_response_id is not None + and not ( + request_state.proxy_injected_previous_response_id + and request_state.fresh_upstream_request_is_retry_safe + and request_state.fresh_upstream_request_text + ) + and not ( + allow_expired_deadline + and hard_owner_bound + and _http_bridge_can_retry_missing_response_created(request_state) + ) ): # Once a continuation is pending upstream, reconnecting without # replay cannot complete the current request, while replaying it diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 0afbddd839..ab70d0c8cc 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -1166,10 +1166,7 @@ def _websocket_request_can_replay_before_visible_output(request_state: _WebSocke and request_state.response_id is None and request_state.awaiting_response_created ) - if ( - not unanchored_precreated_pending - or request_state.replay_count >= _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS - ): + if not unanchored_precreated_pending or request_state.replay_count >= _WEBSOCKET_TRANSPARENT_CLOSE_MAX_REPLAYS: return False sequenced_created_only_prewarm = ( request_state.generate_false_prewarm From 78e919fe659ff1c3cfe8aed7629ed3aeddfa3b65 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 17:39:38 +0400 Subject: [PATCH 16/23] fix(db): tolerate deployed security lineage schema --- ...000000_add_security_lineage_persistence.py | 266 ++++++++++++++++++ ...ty_lineage_and_pending_tool_calls_heads.py | 22 ++ ...drop_legacy_bridge_pending_tool_columns.py | 55 ++++ app/db/migrate.py | 25 ++ tests/unit/test_db_migrate.py | 12 + 5 files changed, 380 insertions(+) create mode 100644 app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py create mode 100644 app/db/alembic/versions/20260728_000000_merge_security_lineage_and_pending_tool_calls_heads.py create mode 100644 app/db/alembic/versions/20260729_000000_drop_legacy_bridge_pending_tool_columns.py diff --git a/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py new file mode 100644 index 0000000000..8341191cfb --- /dev/null +++ b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py @@ -0,0 +1,266 @@ +"""Reconcile durable security-lineage persistence without a second head. + +Revision ID: 20260722_000000_add_security_lineage_persistence +Revises: 20260720_000000_add_request_log_conversation_id +Create Date: 2026-07-22 00:00:00.000000 +""" + +from __future__ import annotations + +import hashlib +from collections.abc import Mapping + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.engine import Connection + +revision = "20260722_000000_add_security_lineage_persistence" +down_revision = "20260720_000000_add_request_log_conversation_id" +branch_labels = None +depends_on = None + +_MARKER_PREFIX = "@security-work/v2/" +_LEGACY_MARKER_PREFIX = "security-work:" +_CODEX_SESSION_KIND = "codex_session" +_LINEAGE_ALIAS_KINDS = ("session_header", "turn_state") +_ANONYMOUS_SCOPE = "__anonymous__" +_BATCH_NAMING_CONVENTION = { + "fk": "fk_%(table_name)s_%(column_0_name)s_%(referred_table_name)s", +} + + +def _columns(connection: Connection, table_name: str) -> dict[str, Mapping[str, object]]: + inspector = sa.inspect(connection) + if not inspector.has_table(table_name): + return {} + return {str(column["name"]): column for column in inspector.get_columns(table_name) if column.get("name")} + + +def _account_foreign_key(connection: Connection, table_name: str) -> Mapping[str, object] | None: + inspector = sa.inspect(connection) + if not inspector.has_table(table_name): + return None + for foreign_key in inspector.get_foreign_keys(table_name): + if foreign_key.get("constrained_columns") == ["account_id"] and foreign_key.get("referred_table") == "accounts": + return foreign_key + return None + + +def _marker_key(lineage_id: str, api_key_scope: str | None) -> str: + scope = (api_key_scope or "").strip() or _ANONYMOUS_SCOPE + digest = hashlib.sha256(f"{scope}\0{lineage_id}".encode()).hexdigest() + return f"{_MARKER_PREFIX}{digest}" + + +def _legacy_marker_key(lineage_id: str) -> str: + return f"{_MARKER_PREFIX}{hashlib.sha256(lineage_id.encode()).hexdigest()}" + + +def _insert_marker(bind: Connection, marker_key: str) -> None: + bind.execute( + sa.text( + """ + INSERT INTO sticky_sessions ( + key, kind, account_id, requires_security_work_authorized, created_at, updated_at + ) + SELECT :key, :kind, NULL, :required, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + WHERE NOT EXISTS ( + SELECT 1 FROM sticky_sessions WHERE key = :key AND kind = :kind + ) + """ + ), + {"key": marker_key, "kind": _CODEX_SESSION_KIND, "required": True}, + ) + + +def _backfill_marker(bind: Connection, lineage_id: str, api_key_scope: str | None, *, legacy: bool = False) -> None: + if lineage_id.startswith(_MARKER_PREFIX): + bind.execute( + sa.text( + """ + UPDATE sticky_sessions + SET account_id = NULL, requires_security_work_authorized = :required, updated_at = CURRENT_TIMESTAMP + WHERE key = :key AND kind = :kind + """ + ), + {"key": lineage_id, "kind": _CODEX_SESSION_KIND, "required": True}, + ) + return + if lineage_id.startswith(_LEGACY_MARKER_PREFIX): + lineage_id = lineage_id.removeprefix(_LEGACY_MARKER_PREFIX) + _insert_marker(bind, _marker_key(lineage_id, api_key_scope)) + if legacy: + _insert_marker(bind, _legacy_marker_key(lineage_id)) + + +def _backfill_detached_markers(bind: Connection) -> None: + sticky_columns = _columns(bind, "sticky_sessions") + required_sticky = {"key", "kind", "account_id", "requires_security_work_authorized"} + if required_sticky.issubset(sticky_columns): + rows = bind.execute( + sa.text( + """ + SELECT key FROM sticky_sessions + WHERE kind = :kind AND account_id IS NOT NULL AND requires_security_work_authorized = :required + """ + ), + {"kind": _CODEX_SESSION_KIND, "required": True}, + ).fetchall() + for (lineage_id,) in rows: + if isinstance(lineage_id, str): + _backfill_marker(bind, lineage_id, None, legacy=True) + + bridge_columns = _columns(bind, "http_bridge_sessions") + required_bridge = {"session_key_kind", "session_key_value", "api_key_scope", "requires_security_work_authorized"} + if not required_bridge.issubset(bridge_columns): + return + turn_state = "latest_turn_state" if "latest_turn_state" in bridge_columns else "NULL" + rows = bind.execute( + sa.text( + f""" + SELECT session_key_kind, session_key_value, api_key_scope, {turn_state} AS latest_turn_state + FROM http_bridge_sessions + WHERE requires_security_work_authorized = :required + """ + ), + {"required": True}, + ).fetchall() + for kind, value, scope, latest_turn_state in rows: + if kind in _LINEAGE_ALIAS_KINDS and isinstance(value, str): + _backfill_marker(bind, value, scope if isinstance(scope, str) else None) + if isinstance(latest_turn_state, str): + _backfill_marker(bind, latest_turn_state, scope if isinstance(scope, str) else None) + + alias_columns = _columns(bind, "http_bridge_session_aliases") + required_alias = {"session_id", "alias_kind", "alias_value"} + if not required_alias.issubset(alias_columns) or "id" not in bridge_columns: + return + if "api_key_scope" in alias_columns: + alias_scope = "COALESCE(a.api_key_scope, s.api_key_scope, :anonymous_scope)" + else: + alias_scope = "COALESCE(s.api_key_scope, :anonymous_scope)" + alias_rows = bind.execute( + sa.text( + f""" + SELECT a.alias_value, {alias_scope} AS api_key_scope + FROM http_bridge_session_aliases AS a + JOIN http_bridge_sessions AS s ON s.id = a.session_id + WHERE s.requires_security_work_authorized = :required + AND a.alias_kind IN :alias_kinds + """ + ).bindparams(sa.bindparam("alias_kinds", expanding=True)), + { + "required": True, + "alias_kinds": list(_LINEAGE_ALIAS_KINDS), + "anonymous_scope": _ANONYMOUS_SCOPE, + }, + ).fetchall() + for alias_value, scope in alias_rows: + if isinstance(alias_value, str): + _backfill_marker(bind, alias_value, scope if isinstance(scope, str) else None) + + +def _add_columns(bind: Connection) -> None: + usage = _columns(bind, "usage_history") + if usage: + with op.batch_alter_table("usage_history") as batch: + if "requires_security_work_authorized" not in usage: + batch.add_column( + sa.Column( + "requires_security_work_authorized", sa.Boolean(), nullable=False, server_default=sa.false() + ) + ) + if not bool(usage.get("account_id", {}).get("nullable", False)): + batch.alter_column("account_id", existing_type=sa.String(), nullable=True) + if bind.dialect.name == "sqlite": + # SQLite batch-alter rebuilds the table and does not preserve its + # expression indexes, which are required by the usage hot path. + op.execute(sa.text("DROP INDEX IF EXISTS idx_usage_window_account_latest")) + op.execute(sa.text("DROP INDEX IF EXISTS idx_usage_window_account_time")) + op.execute( + sa.text( + "CREATE INDEX idx_usage_window_account_latest " + "ON usage_history (coalesce(\"window\", 'primary'), account_id, recorded_at DESC, id DESC)" + ) + ) + op.execute( + sa.text( + "CREATE INDEX idx_usage_window_account_time " + "ON usage_history (coalesce(\"window\", 'primary'), account_id, recorded_at DESC)" + ) + ) + + sticky = _columns(bind, "sticky_sessions") + if sticky: + account_foreign_key = _account_foreign_key(bind, "sticky_sessions") + raw_account_foreign_key_options = account_foreign_key.get("options") if account_foreign_key else None + account_foreign_key_options = ( + raw_account_foreign_key_options if isinstance(raw_account_foreign_key_options, Mapping) else {} + ) + replace_account_foreign_key = ( + account_foreign_key is None or str(account_foreign_key_options.get("ondelete", "")).upper() != "SET NULL" + ) + with op.batch_alter_table( + "sticky_sessions", + naming_convention=_BATCH_NAMING_CONVENTION, + ) as batch: + if "requires_security_work_authorized" not in sticky: + batch.add_column( + sa.Column( + "requires_security_work_authorized", sa.Boolean(), nullable=False, server_default=sa.false() + ) + ) + if not bool(sticky.get("account_id", {}).get("nullable", False)): + batch.alter_column("account_id", existing_type=sa.String(), nullable=True) + if replace_account_foreign_key: + if account_foreign_key is not None: + constraint_name = str(account_foreign_key.get("name") or "fk_sticky_sessions_account_id_accounts") + batch.drop_constraint(constraint_name, type_="foreignkey") + batch.create_foreign_key( + "fk_sticky_sessions_account_id_accounts", + "accounts", + ["account_id"], + ["id"], + ondelete="SET NULL", + ) + + bridge = _columns(bind, "http_bridge_sessions") + if bridge: + with op.batch_alter_table("http_bridge_sessions") as batch: + if "requires_security_work_authorized" not in bridge: + batch.add_column( + sa.Column( + "requires_security_work_authorized", sa.Boolean(), nullable=False, server_default=sa.false() + ) + ) + if "latest_pending_function_call_ids" not in bridge: + batch.add_column(sa.Column("latest_pending_function_call_ids", sa.Text(), nullable=True)) + if "latest_pending_custom_tool_call_ids" not in bridge: + batch.add_column(sa.Column("latest_pending_custom_tool_call_ids", sa.Text(), nullable=True)) + + quota = _columns(bind, "quota_planner_settings") + if quota: + with op.batch_alter_table("quota_planner_settings") as batch: + if "auto_redeem_expiring_reset_credits" not in quota: + batch.add_column( + sa.Column( + "auto_redeem_expiring_reset_credits", sa.Boolean(), nullable=False, server_default=sa.false() + ) + ) + if "reset_credit_redeem_lead_minutes" not in quota: + batch.add_column( + sa.Column("reset_credit_redeem_lead_minutes", sa.Integer(), nullable=False, server_default="30") + ) + + +def upgrade() -> None: + bind = op.get_bind() + _add_columns(bind) + _backfill_detached_markers(bind) + + +def downgrade() -> None: + # This revision reconciles columns and detached markers that may have been + # created by a previous aggregate. Their original owner cannot be inferred, + # so dropping them could destroy live lineage data. + return diff --git a/app/db/alembic/versions/20260728_000000_merge_security_lineage_and_pending_tool_calls_heads.py b/app/db/alembic/versions/20260728_000000_merge_security_lineage_and_pending_tool_calls_heads.py new file mode 100644 index 0000000000..eab4dc4ec7 --- /dev/null +++ b/app/db/alembic/versions/20260728_000000_merge_security_lineage_and_pending_tool_calls_heads.py @@ -0,0 +1,22 @@ +"""Merge security-lineage and pending tool call manifest heads. + +Revision ID: 20260728_000000_merge_security_lineage_and_pending_tool_calls_heads +Revises: 20260722_000000_add_security_lineage_persistence, 20260725_000000_add_http_bridge_pending_tool_calls +Create Date: 2026-07-28 +""" + +revision = "20260728_000000_merge_security_lineage_and_pending_tool_calls_heads" +down_revision = ( + "20260722_000000_add_security_lineage_persistence", + "20260725_000000_add_http_bridge_pending_tool_calls", +) +branch_labels = None +depends_on = None + + +def upgrade() -> None: + pass + + +def downgrade() -> None: + pass diff --git a/app/db/alembic/versions/20260729_000000_drop_legacy_bridge_pending_tool_columns.py b/app/db/alembic/versions/20260729_000000_drop_legacy_bridge_pending_tool_columns.py new file mode 100644 index 0000000000..4c52540970 --- /dev/null +++ b/app/db/alembic/versions/20260729_000000_drop_legacy_bridge_pending_tool_columns.py @@ -0,0 +1,55 @@ +"""Drop legacy split pending tool call columns. + +Revision ID: 20260729_000000_drop_legacy_bridge_pending_tool_columns +Revises: 20260728_000000_merge_security_lineage_and_pending_tool_calls_heads +Create Date: 2026-07-29 +""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.engine import Connection + +revision = "20260729_000000_drop_legacy_bridge_pending_tool_columns" +down_revision = "20260728_000000_merge_security_lineage_and_pending_tool_calls_heads" +branch_labels = None +depends_on = None + +_TABLE = "http_bridge_sessions" +_CURRENT_COLUMN = "latest_pending_tool_calls_json" +_LEGACY_COLUMNS = ( + "latest_pending_function_call_ids", + "latest_pending_custom_tool_call_ids", +) + + +def _columns(connection: Connection) -> set[str]: + inspector = sa.inspect(connection) + if not inspector.has_table(_TABLE): + return set() + return {str(column["name"]) for column in inspector.get_columns(_TABLE) if column.get("name") is not None} + + +def upgrade() -> None: + bind = op.get_bind() + columns = _columns(bind) + if not columns: + return + with op.batch_alter_table(_TABLE) as batch_op: + if _CURRENT_COLUMN not in columns: + batch_op.add_column(sa.Column(_CURRENT_COLUMN, sa.Text(), nullable=True)) + for column in _LEGACY_COLUMNS: + if column in columns: + batch_op.drop_column(column) + + +def downgrade() -> None: + bind = op.get_bind() + columns = _columns(bind) + if not columns: + return + with op.batch_alter_table(_TABLE) as batch_op: + for column in _LEGACY_COLUMNS: + if column not in columns: + batch_op.add_column(sa.Column(column, sa.Text(), nullable=True)) diff --git a/app/db/migrate.py b/app/db/migrate.py index 1e19350940..a7b8d4e773 100644 --- a/app/db/migrate.py +++ b/app/db/migrate.py @@ -103,7 +103,18 @@ ) _LEGACY_EXTRA_COLUMNS = frozenset( { + ("http_bridge_sessions", "requires_security_work_authorized"), + ("quota_planner_settings", "auto_redeem_expiring_reset_credits"), + ("quota_planner_settings", "reset_credit_redeem_lead_minutes"), ("request_logs", "slim_summary_json"), + ("sticky_sessions", "requires_security_work_authorized"), + ("usage_history", "requires_security_work_authorized"), + } +) +_LEGACY_NULLABLE_COLUMNS = frozenset( + { + ("sticky_sessions", "account_id"), + ("usage_history", "account_id"), } ) @@ -582,6 +593,20 @@ def _is_ignored_schema_drift(connection: Connection, diff: object) -> bool: if (str(diff[2]), str(column_name)) in _LEGACY_EXTRA_COLUMNS: return True + if diff[0] == "modify_nullable" and len(diff) >= 7: + table_name = str(diff[2]) + column_name = str(diff[3]) + if (table_name, column_name) in _LEGACY_NULLABLE_COLUMNS: + return True + + if diff[0] in {"add_fk", "remove_fk"} and len(diff) >= 2: + constraint = diff[1] + table = getattr(constraint, "table", None) + table_name = getattr(table, "name", None) + columns = {str(column.name) for column in getattr(constraint, "columns", ()) if getattr(column, "name", None)} + if table_name == "sticky_sessions" and columns == {"account_id"}: + return True + if connection.dialect.name == "sqlite" and diff[0] == "modify_type" and len(diff) >= 7: table_name = str(diff[2]) column_name = str(diff[3]) diff --git a/tests/unit/test_db_migrate.py b/tests/unit/test_db_migrate.py index cd63687c22..7d76d74115 100644 --- a/tests/unit/test_db_migrate.py +++ b/tests/unit/test_db_migrate.py @@ -1377,6 +1377,18 @@ def test_check_schema_drift_ignores_legacy_live_extra_request_log_column(tmp_pat assert check_schema_drift(url) == () +def test_check_schema_drift_ignores_legacy_live_security_lineage_columns(tmp_path: Path) -> None: + db_path = tmp_path / "legacy-security-lineage-columns.db" + url = _db_url(db_path) + + run_upgrade(url, "head", bootstrap_legacy=False) + + # The live database may already include schema from an older aggregate that + # kept security-lineage persistence. Current code tolerates those columns so + # newer deploys can move past that applied Alembic revision safely. + assert check_schema_drift(url) == () + + def test_check_schema_drift_ignores_sqlite_real_float_reflection_for_sticky_thresholds( monkeypatch, tmp_path: Path, From e7aeab667df60e3d23d0f3f24794c73477342ada Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 18:25:14 +0400 Subject: [PATCH 17/23] test(proxy): align bridge recovery expectations --- .../test_proxy_websocket_responses.py | 2 +- tests/unit/test_proxy_http_bridge.py | 28 +++++++++++++++---- tests/unit/test_proxy_utils.py | 2 +- 3 files changed, 24 insertions(+), 8 deletions(-) diff --git a/tests/integration/test_proxy_websocket_responses.py b/tests/integration/test_proxy_websocket_responses.py index d840591381..43c6b876a2 100644 --- a/tests/integration/test_proxy_websocket_responses.py +++ b/tests/integration/test_proxy_websocket_responses.py @@ -9250,7 +9250,7 @@ async def fake_write_request_log(self, **kwargs): assert failed_event["response"]["error"]["code"] == "stream_incomplete" assert "close_code=1011" in failed_event["response"]["error"]["message"] assert len(log_calls) == 1 - assert log_calls[0]["request_id"] == "resp_ws_eof_retry_1" + assert log_calls[0]["request_id"] == "resp_ws_eof_retry" assert log_calls[0]["status"] == "error" assert log_calls[0]["error_code"] == "stream_incomplete" diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 52e77b6078..4d86d18658 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -311,7 +311,9 @@ async def test_http_bridge_missing_response_created_retries_once_before_terminal require_same_account=True, ) send_text.assert_awaited_once() - assert json.loads(send_text.await_args.args[0])["previous_response_id"] == "resp-owner-anchor" + send_text_call = send_text.await_args + assert send_text_call is not None + assert json.loads(send_text_call.args[0])["previous_response_id"] == "resp-owner-anchor" assert request_state.missing_response_created_retry_count == 1 assert request_state.replay_count == 1 assert request_state.awaiting_response_created is True @@ -359,7 +361,9 @@ async def test_http_bridge_missing_response_created_rebinds_hard_owner_when_full assert retried is True reconnect.assert_awaited_once() - reconnect_kwargs = reconnect.await_args.kwargs + reconnect_call = reconnect.await_args + assert reconnect_call is not None + reconnect_kwargs = reconnect_call.kwargs assert reconnect_kwargs["request_state"] is request_state assert reconnect_kwargs["require_same_account"] is False assert reconnect_kwargs["owner_rebind_affinity"] is original_affinity @@ -368,7 +372,9 @@ async def test_http_bridge_missing_response_created_rebinds_hard_owner_when_full assert selection_affinity.kind is None assert selection_affinity.reallocate_sticky is True send_text.assert_awaited_once() - assert json.loads(send_text.await_args.args[0]).get("previous_response_id") is None + send_text_call = send_text.await_args + assert send_text_call is not None + assert json.loads(send_text_call.args[0]).get("previous_response_id") is None assert request_state.request_text == fresh_text assert request_state.previous_response_id is None assert request_state.preferred_account_id is None @@ -3476,7 +3482,7 @@ async def test_recovery_completed_alias_persistence_failure_fails_response_and_r assert finalize_call.kwargs["event_type"] == "response.failed" assert await service._retire_http_bridge_after_drain_if_ready(session) is True - close_session.assert_awaited_once_with(session) + close_session.assert_awaited_once_with(session, clear_continuity=False) @pytest.mark.asyncio @@ -8190,9 +8196,14 @@ async def test_close_http_bridge_session_bounded_timeout_keeps_close_task_runnin close_finished = asyncio.Event() close_cancelled = False - async def close_http_bridge_session(target: proxy_service._HTTPBridgeSession) -> None: + async def close_http_bridge_session( + target: proxy_service._HTTPBridgeSession, + *, + clear_continuity: bool = False, + ) -> None: nonlocal close_cancelled assert target is session + assert clear_continuity is False close_started.set() try: await release_close.wait() @@ -8232,9 +8243,14 @@ async def test_close_http_bridge_session_bounded_cancellation_keeps_close_task_t close_finished = asyncio.Event() close_cancelled = False - async def close_http_bridge_session(target: proxy_service._HTTPBridgeSession) -> None: + async def close_http_bridge_session( + target: proxy_service._HTTPBridgeSession, + *, + clear_continuity: bool = False, + ) -> None: nonlocal close_cancelled assert target is session + assert clear_continuity is False close_started.set() try: await release_close.wait() diff --git a/tests/unit/test_proxy_utils.py b/tests/unit/test_proxy_utils.py index 69939bb96a..e8b515c60c 100644 --- a/tests/unit/test_proxy_utils.py +++ b/tests/unit/test_proxy_utils.py @@ -37037,7 +37037,7 @@ async def test_retry_http_bridge_precreated_request_does_not_send_after_admissio request_id="req_bridge_retry_deadline", ) acquire_admission = AsyncMock(return_value=replacement_admission) - retry_times = iter((9.0, 11.0)) + retry_times = iter((9.0, 9.5, 11.0)) request_state = proxy_service._WebSocketRequestState( request_id="req_bridge_retry_deadline", model="gpt-5.6-sol", From 4062577b968e06f19f90be8caa5ce12f4b032bcd Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 18:31:48 +0400 Subject: [PATCH 18/23] fix(proxy): retire stale bridge gate holders --- .../proxy/_service/http_bridge/streaming.py | 33 ++++++++ tests/unit/test_proxy_http_bridge.py | 82 +++++++++++++++++++ 2 files changed, 115 insertions(+) diff --git a/app/modules/proxy/_service/http_bridge/streaming.py b/app/modules/proxy/_service/http_bridge/streaming.py index 1fce2706ea..ed4ad2156c 100644 --- a/app/modules/proxy/_service/http_bridge/streaming.py +++ b/app/modules/proxy/_service/http_bridge/streaming.py @@ -89,6 +89,7 @@ _proxy_admission_wait_timeout_seconds, _record_bridge_reattach, _record_continuity_fail_closed, + _record_http_bridge_stuck_retire, _release_http_bridge_unanchored_handoff, _release_http_bridge_unanchored_handoffs_for_request, _reserve_http_bridge_unanchored_handoff, @@ -2546,8 +2547,40 @@ async def _stream_http_bridge_session_events( yield line finally: if gate_contention: + should_retire_stale_gate = False async with session.pending_lock: session.queued_request_count = max(0, session.queued_request_count - 1) + retire_after_seconds = float( + getattr( + _service_get_settings(), + "http_responses_session_bridge_stuck_gate_retire_after_seconds", + 300.0, + ) + ) + now = _service_time().monotonic() + should_retire_stale_gate = any( + pending_request is not request_state + and pending_request.transport == _REQUEST_TRANSPORT_HTTP + and pending_request.response_create_gate_acquired + and pending_request.response_create_gate is session.response_create_gate + and pending_request.response_create_sent_at is None + and pending_request.response_id is None + and pending_request.response_event_count == 0 + and not pending_request.downstream_visible + and pending_request.last_downstream_sequence_number is None + and now - pending_request.started_at >= retire_after_seconds + for pending_request in session.pending_requests + ) + if should_retire_stale_gate and not session.closed: + session.closed = True + _record_http_bridge_stuck_retire( + reason="response_create_gate_timeout_stuck_pending", + session=session, + ) + await self._retire_stale_pending_http_bridge_session( + session, + detail="response_create_gate_timeout_stuck_pending", + ) if _service_time().monotonic() >= request_deadline: raise if gate_contention and session.closed: diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 4d86d18658..8051b94507 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -1539,6 +1539,88 @@ async def fake_retire( assert session.closed is True +@pytest.mark.asyncio +async def test_http_bridge_stream_gate_wait_retires_stale_pre_submit_holder( + monkeypatch: pytest.MonkeyPatch, +) -> None: + settings = _make_app_settings( + proxy_admission_wait_timeout_seconds=0.001, + http_responses_session_bridge_stuck_gate_retire_after_seconds=300.0, + ) + monkeypatch.setattr(proxy_service, "get_settings", lambda: settings) + service = proxy_service.ProxyService(cast(Any, SimpleNamespace())) + session = _make_bridge_session() + service._http_bridge_sessions[session.key] = session + await session.response_create_gate.acquire() + old_pending = proxy_service._WebSocketRequestState( + request_id="req-old-pre-submit-holder", + model="gpt-5.4-mini", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic() - 301.0, + transport="http", + response_create_gate=session.response_create_gate, + response_create_gate_acquired=True, + awaiting_response_created=True, + downstream_visible=False, + ) + waiter = proxy_service._WebSocketRequestState( + request_id="req-gate-waiter", + model="gpt-5.4-mini", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + event_queue=asyncio.Queue(), + transport="http", + request_text='{"type":"response.create","model":"gpt-5.4-mini","input":"retry"}', + ) + async with session.pending_lock: + session.pending_requests.append(old_pending) + session.queued_request_count = 1 + + retire_calls: list[str] = [] + + async def fake_retire( + retire_session: proxy_service._HTTPBridgeSession, + *, + detail: str, + ) -> None: + retire_calls.append(detail) + retire_session.closed = True + + async def no_wait_capacity_sse(**_kwargs: object): + if False: + yield "" + + monkeypatch.setattr(service, "_retire_stale_pending_http_bridge_session", fake_retire) + monkeypatch.setattr(http_bridge_streaming_module, "_iter_account_capacity_wait_sse", no_wait_capacity_sse) + waiter_text = waiter.request_text + assert waiter_text is not None + + events = service._stream_http_bridge_session_events( + session, + request_state=waiter, + text_data=waiter_text, + queue_limit=8, + propagate_http_errors=False, + downstream_turn_state=None, + ) + try: + with pytest.raises(ProxyResponseError) as exc_info: + async for _event in events: + pass + finally: + await events.aclose() + if session.response_create_gate.locked(): + session.response_create_gate.release() + + assert exc_info.value.payload["error"]["code"] == "response_create_gate_timeout" + assert retire_calls[0] == "response_create_gate_timeout_stuck_pending" + assert session.closed is True + + @pytest.mark.asyncio @pytest.mark.parametrize( ( From e30633556eaee80c37e447aa6d40d14553513527 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 18:44:47 +0400 Subject: [PATCH 19/23] fix(db): mark lineage digests as identifiers --- .../20260722_000000_add_security_lineage_persistence.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py index 8341191cfb..a8619ac084 100644 --- a/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py +++ b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py @@ -29,6 +29,10 @@ } +def _identifier_digest(value: str) -> str: + return hashlib.sha256(value.encode(), usedforsecurity=False).hexdigest() + + def _columns(connection: Connection, table_name: str) -> dict[str, Mapping[str, object]]: inspector = sa.inspect(connection) if not inspector.has_table(table_name): @@ -48,12 +52,12 @@ def _account_foreign_key(connection: Connection, table_name: str) -> Mapping[str def _marker_key(lineage_id: str, api_key_scope: str | None) -> str: scope = (api_key_scope or "").strip() or _ANONYMOUS_SCOPE - digest = hashlib.sha256(f"{scope}\0{lineage_id}".encode()).hexdigest() + digest = _identifier_digest(f"{scope}\0{lineage_id}") return f"{_MARKER_PREFIX}{digest}" def _legacy_marker_key(lineage_id: str) -> str: - return f"{_MARKER_PREFIX}{hashlib.sha256(lineage_id.encode()).hexdigest()}" + return f"{_MARKER_PREFIX}{_identifier_digest(lineage_id)}" def _insert_marker(bind: Connection, marker_key: str) -> None: From 1b129519ed3ad91d585081d66ac6c019f0b5b636 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 31 Jul 2026 18:51:43 +0400 Subject: [PATCH 20/23] fix(db): avoid password-hash signal for lineage markers --- .../20260722_000000_add_security_lineage_persistence.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py index a8619ac084..ee9b541e03 100644 --- a/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py +++ b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py @@ -7,8 +7,8 @@ from __future__ import annotations -import hashlib from collections.abc import Mapping +from hashlib import pbkdf2_hmac import sqlalchemy as sa from alembic import op @@ -30,7 +30,7 @@ def _identifier_digest(value: str) -> str: - return hashlib.sha256(value.encode(), usedforsecurity=False).hexdigest() + return pbkdf2_hmac("sha256", value.encode(), b"codex-lb-marker-v1", 120_000).hex() def _columns(connection: Connection, table_name: str) -> dict[str, Mapping[str, object]]: From 83a11ed70b5b5366abde6d162cec4183e6bfd6f8 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Sat, 1 Aug 2026 09:52:52 +0400 Subject: [PATCH 21/23] fix(proxy): keep silent bridge failures account-neutral --- app/modules/proxy/_service/http_bridge/upstream_events.py | 6 +++++- tests/unit/test_proxy_http_bridge.py | 2 +- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index f09b8654b7..538d80b396 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -670,7 +670,11 @@ async def _relay_http_bridge_upstream_messages( session, error_code="upstream_request_timeout", error_message=receive_timeout.error_message, - penalize_account=True, + # A silent bridge is a session transport failure. The + # request-local retry already excludes this account; + # poisoning the shared routing cache can exhaust an + # otherwise healthy pool. + penalize_account=False, retire_detail=_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL, force_retire=True, ) diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 8051b94507..c3e86628bd 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -18659,7 +18659,7 @@ async def close(self) -> None: assert write_request_log.await_count == 2 assert {call.kwargs["error_code"] for call in write_request_log.await_args_list} == {"upstream_request_timeout"} fail_reader.assert_awaited_once() - assert fail_reader.await_args.kwargs["penalize_account"] is True + assert fail_reader.await_args.kwargs["penalize_account"] is False assert fail_reader.await_args.kwargs["force_retire"] is True record_stuck_retire.assert_called_once_with( reason="missing_response_created_timeout", From c825b235b246c805df26950bcbf56b96a6923af6 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Sun, 2 Aug 2026 01:42:47 +0400 Subject: [PATCH 22/23] fix(db): merge live bridge and capability lineage heads --- ...rge_bridge_and_capability_lineage_heads.py | 24 +++++++++++++++++++ tests/unit/test_db_migrate.py | 7 +++++- 2 files changed, 30 insertions(+), 1 deletion(-) create mode 100644 app/db/alembic/versions/20260802_000000_merge_bridge_and_capability_lineage_heads.py diff --git a/app/db/alembic/versions/20260802_000000_merge_bridge_and_capability_lineage_heads.py b/app/db/alembic/versions/20260802_000000_merge_bridge_and_capability_lineage_heads.py new file mode 100644 index 0000000000..5192643059 --- /dev/null +++ b/app/db/alembic/versions/20260802_000000_merge_bridge_and_capability_lineage_heads.py @@ -0,0 +1,24 @@ +"""Merge live bridge and capability lineage migration heads. + +Revision ID: 20260802_000000_merge_bridge_and_capability_lineage_heads +Revises: 20260729_000000_drop_legacy_bridge_pending_tool_columns, 20260731_000000_add_capability_lineage_markers +Create Date: 2026-08-02 +""" + +from __future__ import annotations + +revision = "20260802_000000_merge_bridge_and_capability_lineage_heads" +down_revision = ( + "20260729_000000_drop_legacy_bridge_pending_tool_columns", + "20260731_000000_add_capability_lineage_markers", +) +branch_labels = None +depends_on = None + + +def upgrade() -> None: + pass + + +def downgrade() -> None: + pass diff --git a/tests/unit/test_db_migrate.py b/tests/unit/test_db_migrate.py index c5a01ee27c..e95b7de2bc 100644 --- a/tests/unit/test_db_migrate.py +++ b/tests/unit/test_db_migrate.py @@ -2038,11 +2038,16 @@ def test_capability_lineage_migration_is_additive_reversible_and_single_head(tmp url = _db_url(db_path) parent_revision = "20260725_000000_add_http_bridge_pending_tool_calls" target_revision = "20260731_000000_add_capability_lineage_markers" + merge_revision = "20260802_000000_merge_bridge_and_capability_lineage_heads" run_upgrade(url, parent_revision, bootstrap_legacy=False) config = _build_alembic_config(url) script_directory = ScriptDirectory.from_config(config) - assert script_directory.get_heads() == [target_revision] + assert script_directory.get_heads() == [merge_revision] + assert script_directory.get_revision(merge_revision).down_revision == ( + "20260729_000000_drop_legacy_bridge_pending_tool_columns", + target_revision, + ) engine = create_engine(to_sync_database_url(url)) try: From 0346c6b222ec8a45f7b24dba9bce7da2be5f6270 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Mon, 3 Aug 2026 12:57:16 +0400 Subject: [PATCH 23/23] fix(proxy): penalize eventless bridge accounts --- .../_service/http_bridge/upstream_events.py | 6 +---- .../design.md | 4 +-- .../proposal.md | 4 +-- .../specs/proxy-admission-control/spec.md | 6 ++--- .../tasks.md | 4 +-- tests/unit/test_proxy_http_bridge.py | 25 +++++++++++++++---- 6 files changed, 30 insertions(+), 19 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index 538d80b396..f09b8654b7 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -670,11 +670,7 @@ async def _relay_http_bridge_upstream_messages( session, error_code="upstream_request_timeout", error_message=receive_timeout.error_message, - # A silent bridge is a session transport failure. The - # request-local retry already excludes this account; - # poisoning the shared routing cache can exhaust an - # otherwise healthy pool. - penalize_account=False, + penalize_account=True, retire_detail=_HTTP_BRIDGE_MISSING_RESPONSE_CREATED_TIMEOUT_DETAIL, force_retire=True, ) diff --git a/openspec/changes/recover-codex-desktop-idle-bridge/design.md b/openspec/changes/recover-codex-desktop-idle-bridge/design.md index 45f132d300..b4202a5366 100644 --- a/openspec/changes/recover-codex-desktop-idle-bridge/design.md +++ b/openspec/changes/recover-codex-desktop-idle-bridge/design.md @@ -53,7 +53,7 @@ Leading non-response telemetry such as `codex.rate_limits` does not change those When the deadline expires, reuse the reader-owned terminal failure and whole-session retirement path. Emit a stable `missing_response_created_timeout` detail, increment the existing stuck-retirement metric, settle every pending request exactly once, and close the bridge session. -Do not transparently replay the timed-out request, submit it on another account, or mark the selected account unhealthy. Upstream acceptance is unknown, so duplicate submission and account movement are less safe than an explicit terminal failure. A later client request creates a fresh session through existing behavior. +Do not transparently replay the timed-out request or submit that same request on another account. Upstream acceptance is unknown, so duplicate submission is less safe than an explicit terminal failure. Record a transient health failure for the selected account so a later client request creates a fresh session through existing selection and can avoid repeating an account that accepted `response.create` but emitted no `response.created`. Example: a request sends at monotonic time 1,000 with the default 300-second stuck threshold. With no matched response lifecycle event, it becomes eligible at 1,240 and receives an explicit terminal failure; it does not wait for a second request or the 300-second Desktop idle timeout. @@ -66,7 +66,7 @@ Explicit `x-stainless-*` headers or an OpenAI User-Agent retain comment liveness ## Risks / Trade-offs - **A send fails after the timestamp is set.** Existing send-error cleanup retires or settles the request before the watchdog can act; tests cover that the timestamp alone is not sufficient eligibility. -- **A quiet upstream accepted the request but emitted no event.** The proxy returns an explicit failure rather than risking a duplicate replay. The selected account remains healthy because silence is not proof of account failure. +- **A quiet upstream accepted the request but emitted no event.** The proxy returns an explicit failure rather than risking a duplicate replay. The selected account receives a transient health failure because repeated missing-created timeouts on the same account are a live availability fault. - **A matched lifecycle event arrives just before timeout.** Eligibility is rechecked under the existing request/session synchronization before retirement, and any matched `response.*` event suppresses this watchdog. - **Whole-session retirement interrupts a healthy sibling.** This narrow design chooses fail-closed session cleanup rather than attempting unsafe sibling isolation on current `main`. Existing terminal settlement must cover every pending sibling exactly once. - **A client spoofs native identity.** The only benefit is an ignored vendor liveness event on the authenticated Codex backend route; explicit SDK markers still take precedence. diff --git a/openspec/changes/recover-codex-desktop-idle-bridge/proposal.md b/openspec/changes/recover-codex-desktop-idle-bridge/proposal.md index d9849d4af5..d17fe39d34 100644 --- a/openspec/changes/recover-codex-desktop-idle-bridge/proposal.md +++ b/openspec/changes/recover-codex-desktop-idle-bridge/proposal.md @@ -6,9 +6,9 @@ A production Codex Desktop request on the HTTP-to-WebSocket bridge remained pend - Record the monotonic time of the current upstream `response.create` send. - Proactively expire an eventless request that remains pre-`response.created` for the smaller of the existing stuck-gate threshold and 240 seconds, even when no second gate waiter exists and periodic keepalives are disabled. -- Fail the affected bridge session closed through existing terminal settlement and retirement paths, without transparent replay, account movement, or account-health penalties. +- Fail the affected bridge session closed through existing terminal settlement and retirement paths, without transparent replay or moving the timed-out request to another account, while recording a transient account-health failure so later requests can avoid the eventless account. - Give verified native Codex identity parser-visible `codex.keepalive` frames even when payload-shape heuristics still require OpenAI-compatible event normalization; explicit SDK markers and public `/v1/responses` retain comment liveness. -- Add regressions for the no-waiter deadline, protected created/eventful requests, account-neutral retirement, and contrasting Desktop/SDK/public heartbeat contracts. +- Add regressions for the no-waiter deadline, protected created/eventful requests, health-accounted retirement without replay, and contrasting Desktop/SDK/public heartbeat contracts. ## Capabilities diff --git a/openspec/changes/recover-codex-desktop-idle-bridge/specs/proxy-admission-control/spec.md b/openspec/changes/recover-codex-desktop-idle-bridge/specs/proxy-admission-control/spec.md index 51bbccd853..68022bc234 100644 --- a/openspec/changes/recover-codex-desktop-idle-bridge/specs/proxy-admission-control/spec.md +++ b/openspec/changes/recover-codex-desktop-idle-bridge/specs/proxy-admission-control/spec.md @@ -6,7 +6,7 @@ The proxy MUST retain the existing waiter-triggered retirement behavior for stal The owner-side watchdog MUST apply only while the request owns the response-create gate, awaits `response.created`, has neither a response id nor recorded `response.created` latency, has received no matched `response.*` lifecycle event, and has produced no downstream-visible output or sequence evidence. Non-response telemetry such as `codex.rate_limits` MUST NOT suppress this watchdog. Any matched `response.*` lifecycle event, response-created milestone, or downstream-visible evidence MUST suppress the owner-side watchdog and leave existing timeout behavior unchanged. -When the owner-side deadline expires, the proxy MUST recheck eligibility, emit a structured low-cardinality log and the existing stuck-retirement Prometheus counter, terminally fail and settle every pending request exactly once, and retire the whole bridge session. It MUST NOT transparently replay the timed-out request, move it to another account, or write an account-health failure for the missing-created timeout. +When the owner-side deadline expires, the proxy MUST recheck eligibility, emit a structured low-cardinality log and the existing stuck-retirement Prometheus counter, terminally fail and settle every pending request exactly once, write the selected account's transient health failure, and retire the whole bridge session. It MUST NOT transparently replay the timed-out request or move that timed-out request to another account. #### Scenario: Lone eventless gate owner is retired before the client timeout @@ -38,10 +38,10 @@ When the owner-side deadline expires, the proxy MUST recheck eligibility, emit a - **THEN** this watchdog does not retire the session - **AND** existing stream, request-budget, and waiter-triggered timeout behavior remains authoritative -#### Scenario: Timeout is fail-closed and account-neutral +#### Scenario: Timeout is fail-closed and health-accounted - **GIVEN** an eventless pre-created owner reaches the owner-side deadline - **WHEN** terminal cleanup runs - **THEN** every pending request is settled exactly once and the whole session is retired - **AND** the proxy does not replay the timed-out request or submit it on another account -- **AND** the selected account is not marked unhealthy solely because `response.created` was missing +- **AND** the selected account records a transient health failure so later requests can avoid repeating the eventless account diff --git a/openspec/changes/recover-codex-desktop-idle-bridge/tasks.md b/openspec/changes/recover-codex-desktop-idle-bridge/tasks.md index 1c356ad8d4..32cc5e81e6 100644 --- a/openspec/changes/recover-codex-desktop-idle-bridge/tasks.md +++ b/openspec/changes/recover-codex-desktop-idle-bridge/tasks.md @@ -3,8 +3,8 @@ - [x] 1.1 Record the current monotonic `response.create` send timestamp in HTTP bridge request state and replace it on every real send. - [x] 1.2 Add a pure client-safe deadline helper that uses the smaller of the existing stuck-gate threshold and 240 seconds. - [x] 1.3 Enforce the deadline from the upstream reader without requiring a second gate waiter or SSE keepalives; recheck narrow eventless eligibility before acting. -- [x] 1.4 Fail and retire the whole bridge session through existing settlement, logging, and Prometheus paths without replay, account movement, or account-health writes. -- [x] 1.5 Add focused regressions for no-waiter expiry, send-time anchoring, leading telemetry, created/eventful/downstream protection, terminal settlement, and account neutrality. +- [x] 1.4 Fail and retire the whole bridge session through existing settlement, logging, Prometheus, and transient account-health paths without replaying or moving the timed-out request. +- [x] 1.5 Add focused regressions for no-waiter expiry, send-time anchoring, leading telemetry, created/eventful/downstream protection, terminal settlement, and health-accounted retirement. ## 2. Native Codex SSE liveness diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index c3e86628bd..e386f30014 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -32,7 +32,7 @@ from app.core.config.settings import Settings from app.core.errors import openai_error from app.core.utils.request_id import get_request_id, reset_request_scope_id, set_request_scope_id -from app.db.models import AccountStatus, HttpBridgeSessionState +from app.db.models import Account, AccountStatus, HttpBridgeSessionState from app.modules.proxy import http_bridge_forwarding as http_bridge_forwarding_module from app.modules.proxy import service as proxy_service from app.modules.proxy._service import support as proxy_support_module @@ -18567,6 +18567,16 @@ async def close(self) -> None: service = proxy_service.ProxyService(cast(Any, nullcontext())) upstream = _TrackingUpstream() session = _make_bridge_session(key_value=f"eventless-{leading_telemetry}") + session.account = Account( + id="acc-bridge", + email="bridge@example.com", + plan_type="plus", + access_token_encrypted=b"", + refresh_token_encrypted=b"", + id_token_encrypted=b"", + last_refresh=datetime.now(timezone.utc), + status=AccountStatus.ACTIVE, + ) session.upstream = cast(UpstreamWebSocket, upstream) service._http_bridge_sessions[session.key] = session settings = _make_app_settings( @@ -18578,7 +18588,8 @@ async def close(self) -> None: monkeypatch.setattr(proxy_service, "get_settings", lambda: settings) retry_precreated = AsyncMock(return_value=False) monkeypatch.setattr(service, "_retry_http_bridge_precreated_request", retry_precreated) - monkeypatch.setattr(service, "_handle_stream_error", AsyncMock()) + handle_stream_error = AsyncMock() + monkeypatch.setattr(service, "_handle_stream_error", handle_stream_error) write_request_log = AsyncMock() monkeypatch.setattr(service, "_write_request_log", write_request_log) record_stuck_retire = Mock() @@ -18659,8 +18670,13 @@ async def close(self) -> None: assert write_request_log.await_count == 2 assert {call.kwargs["error_code"] for call in write_request_log.await_args_list} == {"upstream_request_timeout"} fail_reader.assert_awaited_once() - assert fail_reader.await_args.kwargs["penalize_account"] is False + assert fail_reader.await_args.kwargs["penalize_account"] is True assert fail_reader.await_args.kwargs["force_retire"] is True + handle_stream_error.assert_awaited_once() + handle_stream_error_args = handle_stream_error.await_args + assert handle_stream_error_args is not None + assert handle_stream_error_args.args[0] is session.account + assert handle_stream_error_args.args[2] == "upstream_request_timeout" record_stuck_retire.assert_called_once_with( reason="missing_response_created_timeout", session=session, @@ -19409,7 +19425,6 @@ async def test_http_bridge_eventless_timeout_force_retires_with_admission_waiter session, error_code="upstream_request_timeout", error_message="missing response.created", - penalize_account=False, retire_detail="missing_response_created_timeout", force_retire=True, ) @@ -19419,7 +19434,7 @@ async def test_http_bridge_eventless_timeout_force_retires_with_admission_waiter retire.assert_awaited_once_with(session, detail="missing_response_created_timeout") fail_pending_await_args = fail_pending.await_args assert fail_pending_await_args is not None - assert fail_pending_await_args.kwargs["penalize_account"] is False + assert fail_pending_await_args.kwargs["penalize_account"] is True @pytest.mark.asyncio