From 7534d61e3fc8d07262528a0670ed00b41362db41 Mon Sep 17 00:00:00 2001 From: hansu650 <2788086371@qq.com> Date: Thu, 13 Aug 2026 13:19:08 +0800 Subject: [PATCH 1/5] fix(proxy): fence clean-close bridge replacements --- .../proxy/_service/http_bridge/mixin.py | 12 ++++- .../.openspec.yaml | 2 + .../proposal.md | 36 +++++++++++++ .../specs/responses-api-compat/spec.md | 39 ++++++++++++++ .../tasks.md | 19 +++++++ .../integration/test_http_responses_bridge.py | 51 +++++++++++++++++-- 6 files changed, 154 insertions(+), 5 deletions(-) create mode 100644 openspec/changes/fence-clean-close-local-replacement/.openspec.yaml create mode 100644 openspec/changes/fence-clean-close-local-replacement/proposal.md create mode 100644 openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md create mode 100644 openspec/changes/fence-clean-close-local-replacement/tasks.md diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index c5b8b4d12..250a7d5cb 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -513,7 +513,17 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: continuity_error: ProxyResponseError | None = None owner_mismatch_error: ProxyResponseError | None = None owner_forward: _HTTPBridgeOwnerForward | None = None - force_durable_takeover = force_durable_takeover_after_detach + # A retiring reader removes its session from the local registry + # before durable release finishes. If the next request reaches + # creation in that window, there is no local object left to set + # ``force_durable_takeover_after_detach`` even though the durable + # row still represents the detached generation. Advance the epoch + # for that fresh local replacement so the late release is fenced. + unrepresented_current_owner = ( + durable_lookup is not None + and durable_lookup.owner_instance_id == settings.http_responses_session_bridge_instance_id + ) + force_durable_takeover = force_durable_takeover_after_detach or unrepresented_current_owner missing_turn_state_alias = False sessions_to_close_before_create: list[_HTTPBridgeSession] = [] session_to_return_after_close: _HTTPBridgeSession | None = None diff --git a/openspec/changes/fence-clean-close-local-replacement/.openspec.yaml b/openspec/changes/fence-clean-close-local-replacement/.openspec.yaml new file mode 100644 index 000000000..b6b2d1f67 --- /dev/null +++ b/openspec/changes/fence-clean-close-local-replacement/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-08-13 diff --git a/openspec/changes/fence-clean-close-local-replacement/proposal.md b/openspec/changes/fence-clean-close-local-replacement/proposal.md new file mode 100644 index 000000000..e3f8780e0 --- /dev/null +++ b/openspec/changes/fence-clean-close-local-replacement/proposal.md @@ -0,0 +1,36 @@ +## Why + +An HTTP Responses bridge can receive `response.completed` and then observe a +clean upstream WebSocket close while the next request for the same soft +affinity key is starting. The retiring session is removed from the local +registry before its durable release completes. A replacement created in that +window can reuse the same durable owner epoch, allowing the retiring session's +late fenced release to clear the replacement's ownership and produce an +intermittent `bridge_instance_mismatch` response on a single instance. + +## What Changes + +- Treat a durable row owned by the current instance without a reusable local + session as a local replacement boundary. +- Advance the durable owner epoch before publishing that fresh local session, + so cleanup from the detached generation is fenced out. +- Add deterministic integration and repository-level regression coverage for + clean-close replacement overlapping a late durable release. + +## Capabilities + +### New Capabilities + +(none) + +### Modified Capabilities + +- `responses-api-compat` + +## Impact + +- **Code:** HTTP bridge session creation and durable-ownership claim behavior. +- **API:** no schema changes; the second request remains successful instead of + intermittently returning HTTP 409 on a single-instance clean-close rollover. +- **Persistence:** replacement claims advance the existing durable owner epoch; + no migration is required. diff --git a/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md new file mode 100644 index 000000000..dbbb47fe8 --- /dev/null +++ b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md @@ -0,0 +1,39 @@ +## ADDED Requirements + +### Requirement: Fresh local bridge replacements fence detached generations + +When an HTTP Responses bridge creates a fresh local session for a durable key +whose active row still names the current instance, and no reusable local +session represents that durable generation, the replacement claim MUST advance +the durable owner epoch before the replacement is published. A release from +the detached local generation MUST be fenced out and MUST NOT clear the +replacement's owner, lease, active state, or continuity anchors. + +Ordinary reuse of a registered local session MUST NOT advance the owner epoch, +and an active row owned by another live instance MUST continue to follow the +existing owner-forwarding or mismatch behavior. + +#### Scenario: clean-close replacement survives a late local release + +- **GIVEN** a local HTTP bridge session has delivered `response.completed` +- **AND** its upstream WebSocket then closes cleanly +- **AND** retirement detaches that session before its durable release completes +- **WHEN** the next request creates a replacement for the same durable key +- **THEN** the replacement claim advances the durable owner epoch +- **AND** the detached generation's late release is fenced out +- **AND** the next request does not receive `bridge_instance_mismatch` + +#### Scenario: registered local reuse keeps its generation + +- **GIVEN** the durable row and a reusable registered local session represent + the same current owner generation +- **WHEN** another compatible request reuses that session +- **THEN** the durable owner epoch is not advanced solely because of reuse + +#### Scenario: a live remote owner remains protected + +- **GIVEN** a durable bridge row has an unexpired lease owned by another live + instance +- **WHEN** this instance cannot reuse a local session for that key +- **THEN** it does not advance or take over the remote owner's epoch without + the existing takeover authorization diff --git a/openspec/changes/fence-clean-close-local-replacement/tasks.md b/openspec/changes/fence-clean-close-local-replacement/tasks.md new file mode 100644 index 000000000..27d5a8f80 --- /dev/null +++ b/openspec/changes/fence-clean-close-local-replacement/tasks.md @@ -0,0 +1,19 @@ +## 1. Durable replacement contract + +- [x] 1.1 Advance the durable owner epoch when a fresh local session replaces + an unrepresented durable generation owned by the current instance. +- [x] 1.2 Preserve ordinary local reuse and cross-instance ownership checks. + +## 2. Regression coverage + +- [x] 2.1 Add a deterministic product-path test that overlaps clean-close + retirement with replacement creation and proves the second request succeeds. +- [x] 2.2 Add focused ownership coverage proving the detached generation's + late release cannot clear the replacement owner. + +## 3. Verification + +- [x] 3.1 Run the focused HTTP bridge integration and durable-session tests + repeatedly. +- [ ] 3.2 Run strict OpenSpec validation and the relevant lint/type/local CI + gates. diff --git a/tests/integration/test_http_responses_bridge.py b/tests/integration/test_http_responses_bridge.py index d4b3a0fda..725f0533e 100644 --- a/tests/integration/test_http_responses_bridge.py +++ b/tests/integration/test_http_responses_bridge.py @@ -7062,7 +7062,11 @@ async def fake_submit_http_bridge_request( @pytest.mark.asyncio -async def test_v1_responses_http_bridge_reconnects_after_clean_upstream_close(async_client, monkeypatch): +async def test_v1_responses_http_bridge_reconnects_after_clean_upstream_close( + async_client, + app_instance, + monkeypatch, +): # The app lifespan registers the process hostname in the durable bridge # ring before this test installs its settings. Keep the test on that same # instance so the startup heartbeat cannot make the reconnect path look @@ -7074,6 +7078,37 @@ async def test_v1_responses_http_bridge_reconnects_after_clean_upstream_close(as second_upstream = _FakeBridgeUpstreamWebSocket() upstreams = [first_upstream, second_upstream] connect_count = 0 + service = get_proxy_service_for_app(app_instance) + release_started = asyncio.Event() + allow_release = asyncio.Event() + release_finished = asyncio.Event() + replacement_claimed = asyncio.Event() + claim_count = 0 + released_lookup = None + original_claim_live_session = service._durable_bridge.claim_live_session + original_release_live_session = service._durable_bridge.release_live_session + + async def overlap_replacement_claim_with_release(*args, **kwargs): + nonlocal claim_count + claim_count += 1 + lookup = await original_claim_live_session(*args, **kwargs) + if claim_count == 2: + replacement_claimed.set() + await _wait_for_event(release_finished) + assert released_lookup is not None + # Model the repository refresh that can observe the concurrent + # release after the replacement commit. The returned owner must + # still be the replacement generation. + return released_lookup + return lookup + + async def pause_retiring_release(*args, **kwargs): + nonlocal released_lookup + release_started.set() + await _wait_for_event(allow_release) + released_lookup = await original_release_live_session(*args, **kwargs) + release_finished.set() + return released_lookup async def fake_select_account_with_budget( self, @@ -7138,6 +7173,8 @@ async def fail_legacy_stream(*args, **kwargs): 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) monkeypatch.setattr(proxy_module, "core_stream_responses", fail_legacy_stream) + monkeypatch.setattr(service._durable_bridge, "claim_live_session", overlap_replacement_claim_with_release) + monkeypatch.setattr(service._durable_bridge, "release_live_session", pause_retiring_release) payload = { "model": "gpt-5.1", @@ -7149,13 +7186,19 @@ async def fail_legacy_stream(*args, **kwargs): "prompt_cache_key": f"http-bridge-reconnect-thread-{account_id}", } first = await asyncio.wait_for(async_client.post("/v1/responses", json=payload), timeout=_TEST_SYNC_TIMEOUT_SECONDS) - second = await asyncio.wait_for( - async_client.post("/v1/responses", json=payload), timeout=_TEST_SYNC_TIMEOUT_SECONDS - ) + await _wait_for_event(release_started) + second_task = asyncio.create_task(async_client.post("/v1/responses", json=payload)) + await _wait_for_event(replacement_claimed) + allow_release.set() + second = await asyncio.wait_for(second_task, timeout=_TEST_SYNC_TIMEOUT_SECONDS) assert first.status_code == 200 assert second.status_code == 200 + assert claim_count == 2 assert connect_count == 2 + assert released_lookup is not None + assert released_lookup.owner_instance_id == socket.gethostname() + assert released_lookup.owner_epoch == 2 @pytest.mark.asyncio From ff71625164b65263bf86c6c222f19b8a0587664a Mon Sep 17 00:00:00 2001 From: hansu650 <2788086371@qq.com> Date: Thu, 13 Aug 2026 13:31:57 +0800 Subject: [PATCH 2/5] fix(proxy): preserve newly remote bridge owners --- .../proxy/_service/http_bridge/mixin.py | 4 +- .../specs/responses-api-compat/spec.md | 9 ++- tests/unit/test_durable_bridge_sessions.py | 54 ++++++++++++++++++ tests/unit/test_proxy_http_bridge.py | 56 +++++++++++++++++++ 4 files changed, 118 insertions(+), 5 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index 250a7d5cb..dd34fb2b0 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -523,7 +523,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: durable_lookup is not None and durable_lookup.owner_instance_id == settings.http_responses_session_bridge_instance_id ) - force_durable_takeover = force_durable_takeover_after_detach or unrepresented_current_owner + force_durable_takeover = force_durable_takeover_after_detach missing_turn_state_alias = False sessions_to_close_before_create: list[_HTTPBridgeSession] = [] session_to_return_after_close: _HTTPBridgeSession | None = None @@ -1544,7 +1544,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: durable_lookup, force=force_durable_takeover, ), - force_owner_epoch_advance=force_durable_takeover, + force_owner_epoch_advance=force_durable_takeover or unrepresented_current_owner, # 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 diff --git a/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md index dbbb47fe8..ee0c4ae8f 100644 --- a/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md +++ b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md @@ -34,6 +34,9 @@ existing owner-forwarding or mismatch behavior. - **GIVEN** a durable bridge row has an unexpired lease owned by another live instance -- **WHEN** this instance cannot reuse a local session for that key -- **THEN** it does not advance or take over the remote owner's epoch without - the existing takeover authorization +- **AND** this instance previously observed its own detached generation before + the remote owner claimed the row +- **WHEN** this instance attempts the local replacement claim +- **THEN** it revalidates the owner under the durable row lock +- **AND** it does not advance or take over the remote owner's epoch without the + existing takeover authorization diff --git a/tests/unit/test_durable_bridge_sessions.py b/tests/unit/test_durable_bridge_sessions.py index d7c3ca1d8..2731ed74b 100644 --- a/tests/unit/test_durable_bridge_sessions.py +++ b/tests/unit/test_durable_bridge_sessions.py @@ -1658,6 +1658,60 @@ async def test_durable_bridge_forced_generation_advance_fences_same_account_stal assert stale_release.state == "active" +@pytest.mark.asyncio +async def test_durable_bridge_forced_local_advance_does_not_take_over_new_remote_owner( + coordinator: DurableBridgeSessionCoordinator, +) -> None: + local = await coordinator.claim_live_session( + session_key_kind="session_header", + session_key_value="sid-local-replacement-race", + api_key_id=None, + instance_id="instance-a", + owner_process_epoch="process-a", + lease_ttl_seconds=60.0, + account_id="acc-1", + model="gpt-5.4", + service_tier=None, + latest_turn_state="http_turn_1", + latest_response_id="resp_1", + allow_takeover=True, + ) + remote = await coordinator.claim_live_session( + session_key_kind="session_header", + session_key_value="sid-local-replacement-race", + api_key_id=None, + instance_id="instance-b", + owner_process_epoch="process-b", + lease_ttl_seconds=60.0, + account_id="acc-1", + model="gpt-5.4", + service_tier=None, + latest_turn_state="http_turn_1", + latest_response_id="resp_1", + allow_takeover=True, + ) + + stale_local_replacement = await coordinator.claim_live_session( + session_key_kind="session_header", + session_key_value="sid-local-replacement-race", + api_key_id=None, + instance_id="instance-a", + owner_process_epoch="process-a", + lease_ttl_seconds=60.0, + account_id="acc-1", + model="gpt-5.4", + service_tier=None, + latest_turn_state="http_turn_1", + latest_response_id="resp_1", + allow_takeover=False, + force_owner_epoch_advance=True, + ) + + assert remote.owner_epoch == local.owner_epoch + 1 + assert stale_local_replacement.owner_instance_id == "instance-b" + assert stale_local_replacement.owner_epoch == remote.owner_epoch + + @pytest.mark.asyncio async def test_durable_bridge_clear_response_anchor_nulls_anchor_fields_but_keeps_turn_state( coordinator: DurableBridgeSessionCoordinator, diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 1ede9d331..ff5b20a79 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -16321,6 +16321,62 @@ async def test_get_or_create_http_bridge_session_recovers_locally_when_stale_own service._ring_membership.resolve_endpoint.assert_awaited_once_with("instance-old") +@pytest.mark.asyncio +async def test_get_or_create_http_bridge_session_fences_unrepresented_local_owner_without_authorizing_takeover( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + key = proxy_service._HTTPBridgeSessionKey("session_header", "sid-clean-close", None) + created_session = _make_bridge_session(key=key, key_value="sid-clean-close") + monkeypatch.setattr(service, "_prune_http_bridge_sessions_locked", Mock(return_value=[])) + monkeypatch.setattr(service, "_create_http_bridge_session", AsyncMock(return_value=created_session)) + claim_durable = AsyncMock() + monkeypatch.setattr(service, "_claim_durable_http_bridge_session", claim_durable) + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) + monkeypatch.setattr( + proxy_service, + "_active_http_bridge_instance_ring", + AsyncMock(return_value=("instance-a", ["instance-a"])), + ) + + resolved = await service._get_or_create_http_bridge_session( + key, + headers={"x-codex-session-id": "sid-clean-close"}, + affinity=proxy_service._AffinityPolicy( + key="sid-clean-close", + kind=proxy_service.StickySessionKind.CODEX_SESSION, + ), + api_key=None, + request_model="gpt-5.2", + idle_ttl_seconds=120.0, + max_sessions=8, + durable_lookup=proxy_service.DurableBridgeLookup( + session_id="durable-clean-close", + canonical_kind="session_header", + canonical_key="sid-clean-close", + api_key_scope="__anonymous__", + account_id="acc-bridge", + owner_instance_id="instance-a", + owner_epoch=1, + lease_expires_at=proxy_service.utcnow() + timedelta(seconds=60), + state=HttpBridgeSessionState.ACTIVE, + latest_turn_state=None, + latest_response_id=None, + ), + ) + + assert resolved is created_session + claim_durable.assert_awaited_once() + await_args = claim_durable.await_args + assert await_args is not None + assert await_args.kwargs["allow_takeover"] is False + assert await_args.kwargs["force_owner_epoch_advance"] is True + + @pytest.mark.asyncio async def test_get_or_create_http_bridge_session_recovers_locally_when_owner_endpoint_missing_but_replay_anchor_exists( monkeypatch: pytest.MonkeyPatch, From af869152eb4ac0152a951bec77809a9000bc4b26 Mon Sep 17 00:00:00 2001 From: hansu650 <2788086371@qq.com> Date: Thu, 13 Aug 2026 17:06:33 +0800 Subject: [PATCH 3/5] fix(proxy): preserve model-transition owner fencing --- .../proxy/_service/http_bridge/mixin.py | 9 ++- .../specs/responses-api-compat/spec.md | 10 ++++ tests/unit/test_proxy_http_bridge.py | 57 +++++++++++++++++++ 3 files changed, 75 insertions(+), 1 deletion(-) diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index dd34fb2b0..288e2a179 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -437,6 +437,10 @@ async def _get_or_create_http_bridge_session( model_transition_rebind = bool( durable_lookup is not None and not _http_bridge_models_compatible(durable_lookup.model, request_model) ) + durable_lookup_owned_by_current_instance = bool( + durable_lookup is not None + and durable_lookup.owner_instance_id == settings.http_responses_session_bridge_instance_id + ) if model_transition_rebind: durable_lookup = None if await _http_bridge_should_wait_for_registration(self, key, settings): @@ -519,7 +523,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: # ``force_durable_takeover_after_detach`` even though the durable # row still represents the detached generation. Advance the epoch # for that fresh local replacement so the late release is fenced. - unrepresented_current_owner = ( + unrepresented_current_owner = durable_lookup_owned_by_current_instance or ( durable_lookup is not None and durable_lookup.owner_instance_id == settings.http_responses_session_bridge_instance_id ) @@ -719,6 +723,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: model_transition_parent_key = key key = fork_key durable_lookup = None + durable_lookup_owned_by_current_instance = False force_durable_takeover_after_detach = False locally_owned_fork_key = _http_bridge_locally_owned_fork_key( fork_key, forwarded_request, forwarded_original_request_unanchored @@ -1440,6 +1445,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: bind_account_neutral_recovery_owner(session) model_transition_parent_key, key = key, fork_key durable_lookup = None + durable_lookup_owned_by_current_instance = False force_durable_takeover_after_detach = False locally_owned_fork_key = _http_bridge_locally_owned_fork_key( fork_key, forwarded_request, forwarded_original_request_unanchored @@ -1544,6 +1550,7 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: durable_lookup, force=force_durable_takeover, ), + and (force_durable_takeover or not unrepresented_current_owner), force_owner_epoch_advance=force_durable_takeover or unrepresented_current_owner, # restart_takeover means recovering a row whose previous # owner is genuinely gone. Every claim now advances the diff --git a/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md index ee0c4ae8f..cb30e6fac 100644 --- a/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md +++ b/openspec/changes/fence-clean-close-local-replacement/specs/responses-api-compat/spec.md @@ -30,6 +30,16 @@ existing owner-forwarding or mismatch behavior. - **WHEN** another compatible request reuses that session - **THEN** the durable owner epoch is not advanced solely because of reuse +#### Scenario: model transition preserves the detached local generation fence + +- **GIVEN** a durable row still names the current instance after its local + generation has detached +- **AND** the replacement request selects a model incompatible with that row +- **WHEN** model filtering discards the durable lookup before session creation +- **THEN** the replacement claim still advances the durable owner epoch +- **AND** a delayed release from the previous model generation cannot clear the + replacement row + #### Scenario: a live remote owner remains protected - **GIVEN** a durable bridge row has an unexpired lease owned by another live diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index ff5b20a79..3c6de2922 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -16377,6 +16377,63 @@ async def test_get_or_create_http_bridge_session_fences_unrepresented_local_owne assert await_args.kwargs["force_owner_epoch_advance"] is True +@pytest.mark.asyncio +async def test_get_or_create_http_bridge_session_preserves_local_owner_fence_across_model_transition( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + key = proxy_service._HTTPBridgeSessionKey("session_header", "sid-model-transition", None) + created_session = _make_bridge_session(key=key, key_value="sid-model-transition") + monkeypatch.setattr(service, "_prune_http_bridge_sessions_locked", Mock(return_value=[])) + monkeypatch.setattr(service, "_create_http_bridge_session", AsyncMock(return_value=created_session)) + claim_durable = AsyncMock() + monkeypatch.setattr(service, "_claim_durable_http_bridge_session", claim_durable) + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) + monkeypatch.setattr( + proxy_service, + "_active_http_bridge_instance_ring", + AsyncMock(return_value=("instance-a", ["instance-a"])), + ) + + resolved = await service._get_or_create_http_bridge_session( + key, + headers={"x-codex-session-id": "sid-model-transition"}, + affinity=proxy_service._AffinityPolicy( + key="sid-model-transition", + kind=proxy_service.StickySessionKind.CODEX_SESSION, + ), + api_key=None, + request_model="gpt-5.2", + idle_ttl_seconds=120.0, + max_sessions=8, + durable_lookup=proxy_service.DurableBridgeLookup( + session_id="durable-model-transition", + canonical_kind="session_header", + canonical_key="sid-model-transition", + api_key_scope="__anonymous__", + account_id="acc-bridge", + owner_instance_id="instance-a", + owner_epoch=1, + lease_expires_at=proxy_service.utcnow() + timedelta(seconds=60), + state=HttpBridgeSessionState.ACTIVE, + latest_turn_state=None, + latest_response_id=None, + model="gpt-5.1", + ), + ) + + assert resolved is created_session + claim_durable.assert_awaited_once() + await_args = claim_durable.await_args + assert await_args is not None + assert await_args.kwargs["allow_takeover"] is False + assert await_args.kwargs["force_owner_epoch_advance"] is True + + @pytest.mark.asyncio async def test_get_or_create_http_bridge_session_recovers_locally_when_owner_endpoint_missing_but_replay_anchor_exists( monkeypatch: pytest.MonkeyPatch, From fada37f1683cb0909ecf6ed60fa59515bea50c6e Mon Sep 17 00:00:00 2001 From: hansu650 <2788086371@qq.com> Date: Thu, 13 Aug 2026 17:33:04 +0800 Subject: [PATCH 4/5] fix(proxy): carry model-transition owner fencing --- .../proxy/_service/http_bridge/mixin.py | 5 +++- .../proxy/_service/http_bridge/streaming.py | 18 +++++++++++++ tests/unit/test_proxy_http_bridge.py | 27 ++++++++++++++++--- 3 files changed, 45 insertions(+), 5 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index 288e2a179..98a551636 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -349,6 +349,7 @@ async def _get_or_create_http_bridge_session( allow_previous_response_recovery_rebind: bool = False, allow_bootstrap_owner_rebind: bool = False, durable_lookup: DurableBridgeLookup | None = None, + durable_model_transition_owned_by_current_instance: bool = False, request_stage: str = "first_turn", preferred_account_id: str | None = None, preferred_account_has_continuity_provenance: bool = False, @@ -382,6 +383,7 @@ async def _get_or_create_http_bridge_session( allow_previous_response_recovery_rebind: bool = False, allow_bootstrap_owner_rebind: bool = False, durable_lookup: DurableBridgeLookup | None = None, + durable_model_transition_owned_by_current_instance: bool = False, request_stage: str = "first_turn", preferred_account_id: str | None = None, preferred_account_has_continuity_provenance: bool = False, @@ -414,6 +416,7 @@ async def _get_or_create_http_bridge_session( allow_previous_response_recovery_rebind: bool = False, allow_bootstrap_owner_rebind: bool = False, durable_lookup: DurableBridgeLookup | None = None, + durable_model_transition_owned_by_current_instance: bool = False, request_stage: str = "first_turn", preferred_account_id: str | None = None, preferred_account_has_continuity_provenance: bool = False, @@ -437,7 +440,7 @@ async def _get_or_create_http_bridge_session( model_transition_rebind = bool( durable_lookup is not None and not _http_bridge_models_compatible(durable_lookup.model, request_model) ) - durable_lookup_owned_by_current_instance = bool( + durable_lookup_owned_by_current_instance = durable_model_transition_owned_by_current_instance or bool( durable_lookup is not None and durable_lookup.owner_instance_id == settings.http_responses_session_bridge_instance_id ) diff --git a/app/modules/proxy/_service/http_bridge/streaming.py b/app/modules/proxy/_service/http_bridge/streaming.py index d7f2d8f59..cebdc2dc3 100644 --- a/app/modules/proxy/_service/http_bridge/streaming.py +++ b/app/modules/proxy/_service/http_bridge/streaming.py @@ -1511,6 +1511,11 @@ def classify_durable_full_resend( and durable_model_transition_lookup.latest_turn_state is not None ) ) + durable_model_transition_owned_by_current_instance = bool( + durable_model_transition_lookup is not None + and durable_model_transition_lookup.owner_instance_id + == _service_get_settings().http_responses_session_bridge_instance_id + ) if durable_model_transition_lookup is not None: _log_http_bridge_event( "model_transition_isolated", @@ -1533,6 +1538,7 @@ def classify_durable_full_resend( bridge_session_key.api_key_id, ) force_local_recovery_creation = True + durable_model_transition_owned_by_current_instance = False durable_lookup = None dead_owner_anchor = False if durable_lookup is not None: @@ -2061,6 +2067,9 @@ def switch_to_account_neutral_replay() -> None: forwarded_affinity_kind=forwarded_affinity_kind, forwarded_affinity_key=forwarded_affinity_key, durable_lookup=durable_lookup, + durable_model_transition_owned_by_current_instance=( + durable_model_transition_owned_by_current_instance + ), request_stage=request_state.request_stage, preferred_account_id=request_state.preferred_account_id, preferred_account_has_continuity_provenance=preferred_account_has_continuity_provenance, @@ -2330,6 +2339,9 @@ def switch_to_account_neutral_replay() -> None: and not owner_forward_fresh_replay ), durable_lookup=durable_lookup, + durable_model_transition_owned_by_current_instance=( + durable_model_transition_owned_by_current_instance + ), request_stage=( request_state.request_stage if owner_forward_fresh_replay @@ -2950,6 +2962,9 @@ async def rollback_pre_dispatch_recovery_claim() -> None: forwarded_request=forwarded_request, forwarded_original_request_unanchored=original_request_unanchored, durable_lookup=durable_lookup, + durable_model_transition_owned_by_current_instance=( + durable_model_transition_owned_by_current_instance + ), request_stage=request_state.request_stage, preferred_account_id=replacement_preferred_account_id, preferred_account_has_continuity_provenance=preferred_account_has_continuity_provenance, @@ -3313,6 +3328,9 @@ async def rollback_pre_dispatch_recovery_claim() -> None: allow_previous_response_recovery_rebind=allow_previous_response_recovery_rebind, session_header_fallback_key=session_header_fallback_key, durable_lookup=durable_lookup, + durable_model_transition_owned_by_current_instance=( + durable_model_transition_owned_by_current_instance + ), request_stage=retry_request_stage, preferred_account_id=retry_preferred_account_id, preferred_account_has_continuity_provenance=preferred_account_has_continuity_provenance, diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 3c6de2922..7e1eb4f37 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -5147,7 +5147,16 @@ def archive_received(message: UpstreamWebSocketMessage) -> None: archive_received=archive_received, ), ) - monkeypatch.setattr(proxy_service, "get_settings", lambda: _make_app_settings()) + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) + monkeypatch.setattr( + http_bridge_streaming_module, + "_service_get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) monkeypatch.setattr(service, "_retry_http_bridge_precreated_request", AsyncMock(return_value=False)) monkeypatch.setattr(service, "_fail_pending_websocket_requests", AsyncMock()) monkeypatch.setattr(service, "_retire_stale_pending_http_bridge_session", AsyncMock()) @@ -10433,7 +10442,7 @@ def test_verified_durable_full_resend_proof_is_sealed_immutable_and_request_boun canonical_key="sid-proof", api_key_scope="__anonymous__", account_id="acc-proof", - owner_instance_id=None, + owner_instance_id="instance-a", owner_epoch=3, lease_expires_at=None, state=HttpBridgeSessionState.CLOSED, @@ -22193,7 +22202,7 @@ async def test_durable_model_transition_preserves_owner_provenance_when_replacin canonical_key="http_turn_model_parent", api_key_scope="__anonymous__", account_id="acc-model-owner", - owner_instance_id=None, + owner_instance_id="instance-a", owner_epoch=1, lease_expires_at=datetime.now(timezone.utc) + timedelta(seconds=60), state=HttpBridgeSessionState.ACTIVE, @@ -22264,7 +22273,16 @@ async def fail_first_session_before_output( ), ), ) - monkeypatch.setattr(proxy_service, "get_settings", lambda: _make_app_settings()) + monkeypatch.setattr( + proxy_service, + "get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) + monkeypatch.setattr( + http_bridge_streaming_module, + "_service_get_settings", + lambda: _make_app_settings(http_responses_session_bridge_instance_id="instance-a"), + ) monkeypatch.setattr(service._durable_bridge, "lookup_request_targets", AsyncMock(return_value=durable_lookup)) monkeypatch.setattr(service, "_resolve_file_account_for_responses", AsyncMock(return_value=None)) monkeypatch.setattr(service, "_get_or_create_http_bridge_session", fake_get_or_create) @@ -22292,6 +22310,7 @@ async def fail_first_session_before_output( assert all(call["previous_response_id"] is None for call in creation_calls) assert all(call["preferred_account_id"] == "acc-model-owner" for call in creation_calls) assert all(call["preferred_account_has_continuity_provenance"] is True for call in creation_calls) + assert all(call["durable_model_transition_owned_by_current_instance"] is True for call in creation_calls) @pytest.mark.asyncio From 6a8086ad834c6dddffa2439708db414227337243 Mon Sep 17 00:00:00 2001 From: Darafei Praliaskouski Date: Fri, 14 Aug 2026 21:27:18 +0400 Subject: [PATCH 5/5] fix(rebase): complete PR 1710 conflict resolution --- app/modules/proxy/_service/http_bridge/mixin.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/app/modules/proxy/_service/http_bridge/mixin.py b/app/modules/proxy/_service/http_bridge/mixin.py index 98a551636..88fa19b24 100644 --- a/app/modules/proxy/_service/http_bridge/mixin.py +++ b/app/modules/proxy/_service/http_bridge/mixin.py @@ -1549,11 +1549,13 @@ def bind_account_neutral_recovery_owner(session: _HTTPBridgeSession) -> None: await _raise_if_http_bridge_creation_superseded(self, key, inflight_future=inflight_future) await self._claim_durable_http_bridge_session( created_session, - allow_takeover=_http_bridge_claim_allows_takeover( - durable_lookup, - force=force_durable_takeover, + allow_takeover=( + _http_bridge_claim_allows_takeover( + durable_lookup, + force=force_durable_takeover, + ) + and (force_durable_takeover or not unrepresented_current_owner) ), - and (force_durable_takeover or not unrepresented_current_owner), force_owner_epoch_advance=force_durable_takeover or unrepresented_current_owner, # restart_takeover means recovering a row whose previous # owner is genuinely gone. Every claim now advances the