Skip to content
Open
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
7 changes: 6 additions & 1 deletion app/modules/api_keys/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -381,7 +381,12 @@ async def commit(self) -> None:
await self._session.commit()

async def update_last_used(self, key_id: str, *, commit: bool = True) -> None:
"""Compatibility touch for maintenance and durability checks."""
"""Compatibility touch for maintenance and durability checks.

Request settlement uses the write-behind coalescer; this explicit
helper remains available for callers that intentionally need a
transaction-local last-used update.
"""
await self._session.execute(update(ApiKey).where(ApiKey.id == key_id).values(last_used_at=utcnow()))
if commit:
await self._session.commit()
Expand Down
44 changes: 42 additions & 2 deletions app/modules/proxy/_service/http_bridge/request_submit.py
Original file line number Diff line number Diff line change
Expand Up @@ -2702,6 +2702,7 @@ async def _retire_stale_pending_http_bridge_session(
) -> None:
async with session.pending_lock:
retired_request_states = list(session.pending_requests)
baseline_completed_response_id = session.last_completed_response_id
if retired_request_count is None:
retired_request_count = sum(
1
Expand Down Expand Up @@ -2752,13 +2753,52 @@ async def _retire_stale_pending_http_bridge_session(
session,
detail=retry_circuit_detail or detail,
)
session.closed = True

async with session.pending_lock:
Comment thread
Komzpa marked this conversation as resolved.
current_response_events_seen = max(
(getattr(request_state, "response_event_count", 0) for request_state in session.pending_requests),
default=0,
)
current_response_created = any(
request_state.response_id is not None or request_state.latency_response_created_ms is not None
for request_state in session.pending_requests
)
completed_response_id = session.last_completed_response_id
current_event_generation = session.last_upstream_event_generation
# The snapshot above avoids awaiting pending_lock while the global
# registry lock is held, preserving bounded cleanup for other sessions.
caller_response_events_seen = response_events_seen or 0
observed_new_completed_response = (
completed_response_id is not None and completed_response_id != baseline_completed_response_id
)
became_healthy_during_suspend = current_response_events_seen > caller_response_events_seen or (
caller_response_events_seen == 0 and (current_response_created or observed_new_completed_response)
)
should_close = False
async with self._http_bridge_lock:
if session.last_upstream_event_generation != current_event_generation:
became_healthy_during_suspend = True
if session.upstream_close_attempted:
session.closed = True
if self._http_bridge_sessions.get(session.key) is session:
self._http_bridge_sessions.pop(session.key, None)
self._unregister_http_bridge_turn_states_locked(session)
self._unregister_http_bridge_previous_response_ids_locked(session)
return
if became_healthy_during_suspend:
# The pending snapshot is only advisory; a terminal response
# may already have left the deque, and a newer close owner may
# have claimed retirement while this task was suspended.
if not session.upstream_close_attempted:
session.closed = False
session.upstream_control.reconnect_requested = False
session.upstream_control.retire_after_drain = False
return
if self._http_bridge_sessions.get(session.key) is session:
session.closed = True
self._http_bridge_sessions.pop(session.key, None)
self._unregister_http_bridge_turn_states_locked(session)
self._unregister_http_bridge_previous_response_ids_locked(session)
async with session.pending_lock:
should_close = not session.upstream_close_attempted
if should_close:
session.upstream_close_attempted = True
Expand Down
1 change: 1 addition & 0 deletions app/modules/proxy/_service/http_bridge/upstream_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -1623,6 +1623,7 @@ async def _process_parsed_http_bridge_upstream_event(
pending_request_count = len(session.pending_requests)

if matched_request_state is not None:
session.last_upstream_event_generation += 1
now = _service_time().monotonic()
if matched_request_state.latency_first_upstream_event_ms is None:
matched_request_state.latency_first_upstream_event_ms = int(
Expand Down
1 change: 1 addition & 0 deletions app/modules/proxy/_service/support.py
Original file line number Diff line number Diff line change
Expand Up @@ -1056,6 +1056,7 @@ class _HTTPBridgeSession:
durable_session_id: str | None = None
durable_owner_epoch: int | None = None
upstream_reader: asyncio.Task[None] | None = None
last_upstream_event_generation: int = 0
last_upstream_close_code: int | None = None
last_upstream_close_generation: int = 0
closed: bool = False
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-08-06
35 changes: 35 additions & 0 deletions openspec/changes/fix-retire-stale-preawait-liveness/design.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
## Context

`_retire_stale_pending_http_bridge_session` awaits retry-circuit persistence before marking the session closed. During that suspension, upstream event handling can update a pending request's response-event count. The current close path never observes that update.

## Goals / Non-Goals

**Goals:**

- Make the close decision use fresh pending-turn liveness after all suspension points.
- Keep registry removal and session close mutually exclusive with concurrent event bookkeeping.
- Retain genuine stale-session retirement and existing retry-circuit behavior.

**Non-Goals:**

- Redesign retry-circuit persistence or bridge event routing.
- Change stale thresholds, retry details, or public response behavior.

## Decisions

- Hold `_http_bridge_lock` while acquiring `session.pending_lock` for the final decision. This keeps registry identity and pending liveness from being observed as a mixed snapshot.
- Compute the final response-event signal from the current pending request states, while retaining the caller's explicit signal as a lower bound for callers that already observed an event. If any current request has a response id or response-created latency, treat it as healthy too.
- Only then set `session.closed`, unregister the session, and claim the upstream close. If liveness is present, return without mutating retirement state.

## Risks / Trade-offs

- [Lock ordering] Nested locks could deadlock if another path acquires them in reverse order. Existing bridge registry cleanup uses bridge-lock-then-pending-lock ordering; keep the new critical section consistent.
- [Completed request removed] A terminal event removed before the final sample is no longer available in pending state; this change does not invent durable session-wide event history.

## Migration Plan

Deploy as a normal application change. Rollback is a revert of the single commit; no data migration is required.

## Open Questions

None.
22 changes: 22 additions & 0 deletions openspec/changes/fix-retire-stale-preawait-liveness/proposal.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
## Why

Stale HTTP bridge retirement currently decides from a pre-await zero-event snapshot. A healthy pending turn can receive its first response event while retry-circuit bookkeeping suspends, then be closed and removed anyway.

## What Changes

- Re-sample pending-turn liveness immediately before stale retirement closes a session.
- Perform the registry identity, pending-state liveness, and close decision under the existing bridge/session locks.
- Preserve retirement for sessions that remain genuinely eventless.
- Add regression and control coverage for both outcomes.

## Capabilities

### New Capabilities

### Modified Capabilities

- `responses-api-compat`: stale HTTP bridge retirement must not kill a turn that became healthy during retry-circuit suspension.

## Impact

The change is limited to `request_submit.py`, HTTP bridge unit coverage, and the Responses API compatibility OpenSpec delta. No public endpoint or schema changes are introduced.
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
## ADDED Requirements

### Requirement: Stale bridge retirement rechecks liveness after suspension

Before closing and unregistering a stale HTTP bridge session, the service MUST re-sample pending request liveness after retry-circuit bookkeeping awaits. A response event, response id, or equivalent response-created signal newly observed after the caller's pre-suspension snapshot MUST prevent stale retirement. A session that remains eventless MUST still be retired.

#### Scenario: First response event arrives during retry-circuit suspension

- **WHEN** stale retirement samples zero response events and then suspends for retry-circuit bookkeeping
- **AND** a pending turn receives its first response event before the close decision
- **THEN** the final decision observes the event under the bridge and pending-state locks
- **AND** the session remains registered, open, and reusable

#### Scenario: Session remains eventless during retry-circuit suspension

- **WHEN** stale retirement samples zero response events and suspends for retry-circuit bookkeeping
- **AND** no pending turn receives a response or response-created signal
- **THEN** the final decision retires and unregisters the session
16 changes: 16 additions & 0 deletions openspec/changes/fix-retire-stale-preawait-liveness/tasks.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
## 1. Regression coverage

- [x] 1.1 Add a deterministic unit regression that injects a response event during retry-circuit suspension and asserts no retirement.
- [x] 1.2 Add a control assertion that an eventless stale session is still retired.
- [x] 1.3 Run the new regression on origin/main and record the expected failure before implementation.

## 2. Implementation

- [x] 2.1 Re-sample response-event and response-created liveness under the bridge/pending locks immediately before the close decision.
- [x] 2.2 Preserve registry cleanup, close idempotence, retry-circuit accounting, and stale-retire logging.

## 3. Verification

- [x] 3.1 Run the targeted unit regression and control.
- [x] 3.2 Run the HTTP bridge unit and integration suites, recording the known reconnect baseline failures.
- [x] 3.3 Validate OpenSpec and commit the focused change without pushing.
150 changes: 150 additions & 0 deletions tests/unit/test_proxy_http_bridge.py
Original file line number Diff line number Diff line change
Expand Up @@ -24103,6 +24103,156 @@ async def test_http_bridge_retirement_does_not_record_midstream_retry_circuit_fa
close.assert_awaited_once_with(session, reason="retire_stale_pending")


@pytest.mark.asyncio
async def test_http_bridge_retirement_rechecks_response_events_after_retry_suspension(
monkeypatch: pytest.MonkeyPatch,
) -> None:
service = proxy_service.ProxyService(cast(Any, nullcontext()))
request_state = proxy_service._WebSocketRequestState(
request_id="req-first-response-during-retire",
model="gpt-5.6-sol",
service_tier=None,
reasoning_effort=None,
api_key_reservation=None,
started_at=time.monotonic(),
transport="http",
)
session = _make_bridge_session(
key_value="bridge-retire-post-await-liveness",
pending_requests=deque([request_state]),
queued_request_count=1,
)
session.closed = True
session.upstream_control.reconnect_requested = True
session.upstream_control.retire_after_drain = True
service._http_bridge_sessions[session.key] = session
close = AsyncMock()
monkeypatch.setattr(service, "_close_http_bridge_session_bounded", close)

async def record_failure_during_await(*_args: Any, **_kwargs: Any) -> None:
# Model the upstream reader recording response.created while retry
# circuit persistence has suspended the retire coroutine.
await asyncio.sleep(0)
async with session.pending_lock:
request_state.response_event_count = 1
session.last_upstream_event_generation += 1

monkeypatch.setattr(service, "_record_http_bridge_retry_circuit_failure", record_failure_during_await)

await service._retire_stale_pending_http_bridge_session(
session,
detail="stream_incomplete",
response_events_seen=0,
)

assert session.closed is False
assert session.upstream_control.reconnect_requested is False
assert session.upstream_control.retire_after_drain is False
assert service._http_bridge_sessions[session.key] is session
close.assert_not_awaited()


@pytest.mark.asyncio
async def test_http_bridge_retirement_ignores_stale_completed_response_id(
monkeypatch: pytest.MonkeyPatch,
) -> None:
service = proxy_service.ProxyService(cast(Any, nullcontext()))
session = _make_bridge_session(key_value="bridge-retire-stale-completed-id")
session.last_completed_response_id = "resp-old"
service._http_bridge_sessions[session.key] = session
close = AsyncMock()
monkeypatch.setattr(service, "_close_http_bridge_session_bounded", close)
monkeypatch.setattr(service, "_record_http_bridge_retry_circuit_failure", AsyncMock())

await service._retire_stale_pending_http_bridge_session(
session,
detail="stream_incomplete",
response_events_seen=0,
)

assert session.closed is True
assert service._http_bridge_sessions.get(session.key) is None
close.assert_awaited_once_with(session, reason="retire_stale_pending")


@pytest.mark.asyncio
async def test_http_bridge_retirement_rechecks_events_after_pending_snapshot(
monkeypatch: pytest.MonkeyPatch,
) -> None:
service = proxy_service.ProxyService(cast(Any, nullcontext()))
request_state = proxy_service._WebSocketRequestState(
request_id="req-event-after-pending-snapshot",
model="gpt-5.6-sol",
service_tier=None,
reasoning_effort=None,
api_key_reservation=None,
started_at=time.monotonic(),
transport="http",
)
session = _make_bridge_session(
key_value="bridge-retire-post-snapshot-liveness",
pending_requests=deque([request_state]),
queued_request_count=1,
)
session.closed = True
session.upstream_control.reconnect_requested = True
session.upstream_control.retire_after_drain = True
service._http_bridge_sessions[session.key] = session
close = AsyncMock()
monkeypatch.setattr(service, "_close_http_bridge_session_bounded", close)
monkeypatch.setattr(service, "_record_http_bridge_retry_circuit_failure", AsyncMock())

await service._http_bridge_lock.acquire()
try:
retire_task = asyncio.create_task(
service._retire_stale_pending_http_bridge_session(
session,
detail="stream_incomplete",
response_events_seen=0,
)
)
await asyncio.sleep(0)
async with session.pending_lock:
request_state.response_event_count = 1
session.last_upstream_event_generation += 1
finally:
service._http_bridge_lock.release()

await retire_task

assert session.closed is False
assert session.upstream_control.reconnect_requested is False
assert session.upstream_control.retire_after_drain is False
assert service._http_bridge_sessions[session.key] is session
close.assert_not_awaited()


@pytest.mark.asyncio
async def test_http_bridge_retirement_still_closes_eventless_session_after_retry_suspension(
monkeypatch: pytest.MonkeyPatch,
) -> None:
service = proxy_service.ProxyService(cast(Any, nullcontext()))
session = _make_bridge_session(key_value="bridge-retire-eventless-control")
service._http_bridge_sessions[session.key] = session
close = AsyncMock()
monkeypatch.setattr(service, "_close_http_bridge_session_bounded", close)

async def record_failure_during_await(*_args: Any, **_kwargs: Any) -> None:
await asyncio.sleep(0)

monkeypatch.setattr(service, "_record_http_bridge_retry_circuit_failure", record_failure_during_await)

await service._retire_stale_pending_http_bridge_session(
session,
detail="stream_incomplete",
response_events_seen=0,
)

assert session.closed is True
assert service._http_bridge_sessions.get(session.key) is None
close.assert_awaited_once_with(session, reason="retire_stale_pending")


@pytest.mark.asyncio
async def test_http_bridge_retry_circuit_backoff_is_scoped_to_repeated_hard_keys() -> None:
service = proxy_service.ProxyService(cast(Any, nullcontext()))
Expand Down
Loading