Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 14 additions & 9 deletions app/modules/proxy/_service/http_bridge/mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,7 @@ async def _get_or_create_http_bridge_session(
exclude_account_ids: Collection[str] | None = None,
deferred_account_backoff_lifecycle: _DeferredAccountBackoffLifecycle | None = None,
defer_account_health_writes: bool = False,
force_new_turn_state_session: bool = False,
) -> "_HTTPBridgeSession": ...
@overload
async def _get_or_create_http_bridge_session(
Expand Down Expand Up @@ -397,6 +398,7 @@ async def _get_or_create_http_bridge_session(
exclude_account_ids: Collection[str] | None = None,
deferred_account_backoff_lifecycle: _DeferredAccountBackoffLifecycle | None = None,
defer_account_health_writes: bool = False,
force_new_turn_state_session: bool = False,
) -> "_HTTPBridgeSession | _HTTPBridgeOwnerForward": ...
async def _get_or_create_http_bridge_session(
self,
Expand Down Expand Up @@ -429,6 +431,7 @@ async def _get_or_create_http_bridge_session(
exclude_account_ids: Collection[str] | None = None,
deferred_account_backoff_lifecycle: _DeferredAccountBackoffLifecycle | None = None,
defer_account_health_writes: bool = False,
force_new_turn_state_session: bool = False,
) -> "_HTTPBridgeSession | _HTTPBridgeOwnerForward":
settings = _service_get_settings()
request_scope_id = ensure_request_scope_id()
Expand Down Expand Up @@ -678,14 +681,16 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None:
key=key.affinity_key,
):
key = _HTTPBridgeSessionKey("turn_state_header", incoming_turn_state, api_key_id)
elif (
fallback_key := _alias_fallback_key(incoming_session_key, initial_session_key, api_key_id)
) is not None:
key = fallback_key
used_session_header_fallback = True
else:
key = _HTTPBridgeSessionKey("turn_state_header", incoming_turn_state, api_key_id)
missing_turn_state_alias = True
fallback_key = (
_alias_fallback_key(incoming_session_key, initial_session_key, api_key_id)
if not force_new_turn_state_session
else None
)
key = fallback_key or _HTTPBridgeSessionKey(
"turn_state_header", incoming_turn_state, api_key_id
)
missing_turn_state_alias = not (used_session_header_fallback := fallback_key is not None)
pruned_sessions = self._prune_http_bridge_sessions_locked()
if pruned_sessions:
if any(session.key == key for session in pruned_sessions):
Expand Down Expand Up @@ -1269,6 +1274,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None:
incoming_turn_state is not None
and incoming_turn_state.startswith("http_turn_")
and not allow_forward_to_owner
and not force_new_turn_state_session
):
_record_continuity_fail_closed(
surface="http_bridge",
Expand Down Expand Up @@ -1560,8 +1566,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None:
# restart_takeover means recovering a row whose previous
# owner is genuinely gone. Every claim now advances the
# epoch, so epoch > 1 alone would also count ordinary
# local successor claims (no pre-claim lookup, or a
# forced replace of a live local session).
# local successor claims (no pre-claim lookup, or a forced replace of a live local session).
claim_kwargs["record_restart_takeover"] = True
await self._claim_durable_http_bridge_session(created_session, **claim_kwargs)
async with self._http_bridge_lock:
Expand Down
24 changes: 19 additions & 5 deletions app/modules/proxy/_service/http_bridge/streaming.py
Original file line number Diff line number Diff line change
Expand Up @@ -893,6 +893,7 @@ def stream_http_responses(
enforce_openai_sdk_contract: bool = True,
capacity_startup_wait_event: asyncio.Event | None = None,
capacity_startup_ready_event: asyncio.Event | None = None,
verified_v1_goal_restart: bool = False,
) -> AsyncIterator[str]:
_maybe_log_proxy_request_payload("stream_http", payload, headers)
proxy_api_authorization = _header_value_case_insensitive(headers, "authorization")
Expand All @@ -918,6 +919,7 @@ def stream_http_responses(
enforce_openai_sdk_contract=enforce_openai_sdk_contract,
capacity_startup_wait_event=capacity_startup_wait_event,
capacity_startup_ready_event=capacity_startup_ready_event,
verified_v1_goal_restart=verified_v1_goal_restart,
)

async def _stream_http_bridge_or_retry(
Expand All @@ -943,6 +945,7 @@ async def _stream_http_bridge_or_retry(
enforce_openai_sdk_contract: bool = True,
capacity_startup_wait_event: asyncio.Event | None = None,
capacity_startup_ready_event: asyncio.Event | None = None,
verified_v1_goal_restart: bool = False,
) -> AsyncIterator[str]:
dashboard_settings = await _service_get_settings_cache().get()
runtime_config = _http_bridge_runtime_config(dashboard_settings, _service_get_settings())
Expand Down Expand Up @@ -1028,6 +1031,7 @@ async def _stream_http_bridge_or_retry(
enforce_openai_sdk_contract=enforce_openai_sdk_contract,
capacity_startup_wait_event=capacity_startup_wait_event,
capacity_startup_ready_event=capacity_startup_ready_event,
verified_v1_goal_restart=verified_v1_goal_restart,
deferred_account_backoff_tracker=deferred_account_backoff_tracker,
):
yield line
Expand Down Expand Up @@ -1091,6 +1095,7 @@ async def _stream_via_http_bridge(
enforce_openai_sdk_contract: bool = True,
capacity_startup_wait_event: asyncio.Event | None = None,
capacity_startup_ready_event: asyncio.Event | None = None,
verified_v1_goal_restart: bool = False,
deferred_account_backoff_tracker: _DeferredAccountBackoffTracker | None = None,
) -> AsyncIterator[str]:
del suppress_text_done_events
Expand Down Expand Up @@ -2063,6 +2068,7 @@ def switch_to_account_neutral_replay(
exclude_account_ids=fresh_replay_excluded_account_ids or None,
deferred_account_backoff_lifecycle=request_state.deferred_account_backoff_lifecycle,
defer_account_health_writes=request_state.api_key_reservation is not None,
force_new_turn_state_session=verified_v1_goal_restart,
)
except ProxyResponseError as exc:
if not owner_unavailable_allows_account_neutral_replay(exc):
Expand Down Expand Up @@ -2337,6 +2343,7 @@ def switch_to_account_neutral_replay(
exclude_account_ids=request_state.excluded_account_ids or None,
deferred_account_backoff_lifecycle=request_state.deferred_account_backoff_lifecycle,
defer_account_health_writes=request_state.api_key_reservation is not None,
force_new_turn_state_session=verified_v1_goal_restart,
)
except ProxyResponseError as capacity_exc:
if owner_unavailable_allows_account_neutral_replay(capacity_exc):
Expand Down Expand Up @@ -2994,6 +3001,7 @@ def capture_verified_stale_anchor_quarantine_generation(
exclude_account_ids=request_state.excluded_account_ids or None,
deferred_account_backoff_lifecycle=request_state.deferred_account_backoff_lifecycle,
defer_account_health_writes=request_state.api_key_reservation is not None,
force_new_turn_state_session=verified_v1_goal_restart,
)
except ProxyResponseError as capacity_exc:
wait_plan = _http_bridge_capacity_wait_plan(capacity_exc, request_deadline=request_deadline)
Expand Down Expand Up @@ -3100,6 +3108,7 @@ def capture_verified_stale_anchor_quarantine_generation(
exclude_account_ids=request_state.excluded_account_ids or None,
deferred_account_backoff_lifecycle=request_state.deferred_account_backoff_lifecycle,
defer_account_health_writes=request_state.api_key_reservation is not None,
force_new_turn_state_session=verified_v1_goal_restart,
)
except ProxyResponseError as capacity_exc:
wait_plan = _http_bridge_capacity_wait_plan(capacity_exc, request_deadline=request_deadline)
Expand Down Expand Up @@ -3402,6 +3411,7 @@ def capture_verified_stale_anchor_quarantine_generation(
exclude_account_ids=request_state.excluded_account_ids or None,
deferred_account_backoff_lifecycle=request_state.deferred_account_backoff_lifecycle,
defer_account_health_writes=request_state.api_key_reservation is not None,
force_new_turn_state_session=verified_v1_goal_restart,
)
except ProxyResponseError as capacity_exc:
wait_plan = _http_bridge_capacity_wait_plan(capacity_exc, request_deadline=request_deadline)
Expand Down Expand Up @@ -4524,18 +4534,22 @@ def stream_idle_keepalive(*, downstream_response_id: str) -> str | None:
and request_state.error_http_status_override is not None
and request_state.error_http_status_override >= 400
):
if request_state.previous_response_not_found_rewritten:
if request_state.previous_response_not_found_rewritten and not (
request_state.proxy_injected_previous_response_id
and _http_bridge_continuity_bound_without_safe_replay(request_state)
):
raise ProxyResponseError(
request_state.error_http_status_override,
openai_error(
"bridge_previous_response_not_found",
"Upstream websocket closed before response.completed",
),
)
raise ProxyResponseError(
request_state.error_http_status_override,
_openai_error_envelope_from_response_failed_payload(block_payload),
)
if not request_state.previous_response_not_found_rewritten:
raise ProxyResponseError(
request_state.error_http_status_override,
_openai_error_envelope_from_response_failed_payload(block_payload),
)
yield event_block
yielded_any = True
finally:
Expand Down
36 changes: 36 additions & 0 deletions app/modules/proxy/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -698,6 +698,36 @@ def _has_explicit_openai_sdk_marker(request: Request) -> bool:
return "openai" in user_agent


def _verified_v1_goal_restart_headers(
request: Request,
payload: ResponsesRequest,
) -> dict[str, str] | None:
"""Rotate a proved native Goal restart away from its echoed HTTP turn."""
if payload.stream is not True:
return None
if not _is_native_codex_request(request.headers) or _has_explicit_openai_sdk_marker(request):
return None
if proxy_affinity_module._codex_backend_identity(request.headers).thread_id is None:
return None
turn_state = proxy_affinity_module._sticky_key_from_turn_state_header(request.headers)
if (
turn_state is None
or not turn_state.startswith("http_turn_")
or not proxy_affinity_module._is_synthesized_turn_state(turn_state)
):
return None
if not proxy_affinity_module._request_allows_unavailable_legacy_owner_abandonment(payload):
return None

forwarded_headers = dict(request.headers)
forwarded_headers["x-codex-turn-state"] = proxy_affinity_module.ensure_http_downstream_turn_state({})
logger.info(
"v1_goal_restart_turn_state_rotated request_id=%s",
ensure_request_id(),
)
return forwarded_headers


def _is_openai_sdk_request(
request: Request,
payload: V1ResponsesRequest | Mapping[str, JsonValue] | None = None,
Expand Down Expand Up @@ -1303,6 +1333,7 @@ async def v1_responses(
service_tier_was_enforced=service_tier_was_enforced,
)
if responses_payload.stream:
goal_restart_headers = _verified_v1_goal_restart_headers(request, responses_payload)
response = await _stream_responses(
request,
responses_payload,
Expand All @@ -1313,6 +1344,8 @@ async def v1_responses(
prefer_http_bridge=True,
api_key_policy_already_applied=True,
prohibit_fast_mode=prohibit_fast_mode,
forwarded_headers=goal_restart_headers,
verified_v1_goal_restart=goal_restart_headers is not None,
)
else:
response = await _collect_responses(
Expand Down Expand Up @@ -5541,6 +5574,7 @@ async def _stream_responses(
native_codex_heartbeat: bool = False,
api_key_policy_already_applied: bool = False,
prohibit_fast_mode: bool = False,
verified_v1_goal_restart: bool = False,
) -> Response:
# Owner-forwarded payloads have already passed API-key enforcement,
# account-catalog fallback, reservation, and signing on the origin
Expand Down Expand Up @@ -5779,6 +5813,7 @@ def build_response_stream() -> AsyncIterator[str]:
enforce_openai_sdk_contract=enforce_openai_sdk_contract,
capacity_startup_wait_event=capacity_wait_event,
capacity_startup_ready_event=capacity_ready_event,
verified_v1_goal_restart=verified_v1_goal_restart,
)
return context.service.stream_responses(
payload,
Expand Down Expand Up @@ -5832,6 +5867,7 @@ async def _retry() -> AsyncIterator[str]:
enforce_openai_sdk_contract=enforce_openai_sdk_contract,
capacity_startup_wait_event=capacity_wait_event,
capacity_startup_ready_event=capacity_ready_event,
verified_v1_goal_restart=verified_v1_goal_restart,
)
async for line in retry_stream:
yield line
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-08-28
Loading
Loading