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..ee9b541e03 --- /dev/null +++ b/app/db/alembic/versions/20260722_000000_add_security_lineage_persistence.py @@ -0,0 +1,270 @@ +"""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 + +from collections.abc import Mapping +from hashlib import pbkdf2_hmac + +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 _identifier_digest(value: str) -> str: + 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]]: + 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 = _identifier_digest(f"{scope}\0{lineage_id}") + return f"{_MARKER_PREFIX}{digest}" + + +def _legacy_marker_key(lineage_id: str) -> str: + return f"{_MARKER_PREFIX}{_identifier_digest(lineage_id)}" + + +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/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/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/app/modules/proxy/_service/http_bridge/helpers.py b/app/modules/proxy/_service/http_bridge/helpers.py index ca98a7fa9f..b2e7494944 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") @@ -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 0319e918ea..7d64b8de0f 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 @@ -127,6 +128,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, @@ -192,6 +194,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 " @@ -199,6 +202,29 @@ ) +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 + ) + + +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 _rollback_http_bridge_recovery_turn_state_registration( service: Any, receipt: DurableBridgeAliasRegistrationReceipt, @@ -1425,6 +1451,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: @@ -1436,7 +1463,22 @@ 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=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, @@ -1530,12 +1572,18 @@ 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() + 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 ( @@ -1543,7 +1591,19 @@ 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_retry_missing_response_created(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: @@ -1551,15 +1611,35 @@ 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_retry_missing_response_created(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 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 + 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 @@ -1567,6 +1647,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, @@ -1604,7 +1689,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 @@ -1616,6 +1719,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 @@ -1632,7 +1737,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/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/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index a70367af9b..f09b8654b7 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,11 +643,34 @@ 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, + allow_expired_deadline=True, + ) + 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", 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/app/modules/proxy/_service/streaming/mixin.py b/app/modules/proxy/_service/streaming/mixin.py index 6189ae7cd5..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, ) @@ -491,6 +492,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 +579,8 @@ async def _stream_once( settlement.record_success = False settlement.account_health_error = True settlement.error = {"message": error_message} + 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( error_code, diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 141b2b5570..a8a2e020af 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -67,6 +67,8 @@ _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( { "turn_state_header", @@ -786,6 +788,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 @@ -1159,7 +1162,13 @@ 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 @@ -1168,6 +1177,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/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/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/integration/test_http_responses_bridge.py b/tests/integration/test_http_responses_bridge.py index 989783c5af..b12770af52 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, + ) + 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", + "http-bridge-missing-created-retry@example.com", + ) + account = await _get_account(account_id) + silent_upstreams = [_SilentUpstreamWebSocket() for _ in range(5)] + recovered_upstream = _FakeBridgeUpstreamWebSocket() + upstreams = [*silent_upstreams, 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 == 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 + + @pytest.mark.asyncio async def test_backend_responses_http_bridge_retries_precreated_server_overload(async_client, monkeypatch): _install_bridge_settings(monkeypatch, enabled=True) diff --git a/tests/unit/test_db_migrate.py b/tests/unit/test_db_migrate.py index 03a5fc23ed..e95b7de2bc 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, @@ -2026,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: 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 00dd232e8e..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 @@ -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 ) @@ -268,6 +268,123 @@ 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() + 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 + + +@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_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 + 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() + 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 + 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, @@ -1422,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( ( @@ -3365,7 +3564,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 @@ -8079,9 +8278,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() @@ -8121,9 +8325,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() @@ -15306,6 +15515,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 @@ -15371,6 +15581,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 @@ -18356,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( @@ -18367,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() @@ -18444,12 +18666,17 @@ 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() - 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, @@ -18882,12 +19109,76 @@ 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_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, @@ -19134,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, ) @@ -19144,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 @@ -19523,6 +19813,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() @@ -19604,6 +19895,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() diff --git a/tests/unit/test_proxy_utils.py b/tests/unit/test_proxy_utils.py index f5d77f1a88..82d02fcfb0 100644 --- a/tests/unit/test_proxy_utils.py +++ b/tests/unit/test_proxy_utils.py @@ -30122,6 +30122,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() @@ -37083,7 +37144,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",