test(proxy): cover the retry credential blocking the proxy event loop
Extends the off-loop suite to the third and last blocking method on the
`UpstreamAdapter` contract, the 401/429 rotation.
As with the two existing pairs, the primary assertion is **thread identity**,
not latency: a latency assertion measured by an HTTP client on the blocked
loop is vacuous, because the client's own timer cannot advance until the
block ends and it therefore reports a fast response on provably frozen code.
* `test_get_retry_credential_runs_off_the_event_loop` records
`threading.get_ident()` inside the fake adapter and compares it to the
loop thread, and checks the rotation still works end to end (rejected
bearer forwarded first, rotated bearer second).
* `test_event_loop_keeps_running_while_the_retry_credential_resolves`
samples a loop-side heartbeat counter from inside the stalled adapter. On
the unfixed handler it records exactly 0 loop iterations across a 0.5s
rotation.
* `test_retry_credential_failure_still_returns_the_upstream_rejection`
guards the error contract the change must leave alone: a raising rotation
is still swallowed and the upstream's own 401 is streamed back, with no
second forward.
A new `_build_rejecting_upstream` harness drives the `status in {401, 429}`
branch by rejecting every bearer except the rotated one.
This commit is contained in:
@@ -9,7 +9,12 @@ blocking I/O:
|
||||
(``hermes_cli/auth.py:110``), reads ``auth.json`` from disk, and may issue a
|
||||
token-refresh POST. Its terminal-error path takes that lock a second time to
|
||||
persist the quarantined state.
|
||||
* ``NousPortalAdapter.get_retry_credential`` routes to that same
|
||||
``_get_credential`` with ``force_refresh=True``, so the refresh POST it only
|
||||
*may* perform above is unconditional here.
|
||||
* ``XAIGrokAdapter`` reads its key pool off disk under a ``threading.Lock``.
|
||||
Its ``get_retry_credential`` loads the pool and calls
|
||||
``try_refresh_current`` / ``mark_exhausted_and_rotate`` under that lock.
|
||||
|
||||
``create_app`` registers two ``async def`` handlers, so calling those methods
|
||||
directly from a handler freezes the proxy's single event loop — and with it
|
||||
@@ -74,15 +79,22 @@ class _RecordingAdapter(UpstreamAdapter):
|
||||
stall: float = 0.0,
|
||||
ticks: Optional[List[int]] = None,
|
||||
raise_on_credential: bool = False,
|
||||
retry_bearer: Optional[str] = None,
|
||||
raise_on_retry: bool = False,
|
||||
) -> None:
|
||||
self._base_url = base_url
|
||||
self._stall = stall
|
||||
self._ticks = ticks if ticks is not None else [0]
|
||||
self._raise_on_credential = raise_on_credential
|
||||
self._retry_bearer = retry_bearer
|
||||
self._raise_on_retry = raise_on_retry
|
||||
self.credential_thread: Optional[int] = None
|
||||
self.authenticated_thread: Optional[int] = None
|
||||
self.retry_thread: Optional[int] = None
|
||||
self.retry_status_code: Optional[int] = None
|
||||
self.ticks_across_credential: Optional[int] = None
|
||||
self.ticks_across_is_authenticated: Optional[int] = None
|
||||
self.ticks_across_retry: Optional[int] = None
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
@@ -119,8 +131,22 @@ class _RecordingAdapter(UpstreamAdapter):
|
||||
)
|
||||
|
||||
def get_retry_credential(self, *, failed_credential, status_code):
|
||||
_ = failed_credential, status_code
|
||||
return None
|
||||
_ = failed_credential
|
||||
self.retry_thread = threading.get_ident()
|
||||
self.retry_status_code = status_code
|
||||
before = self._ticks[0]
|
||||
if self._stall:
|
||||
time.sleep(self._stall)
|
||||
self.ticks_across_retry = self._ticks[0] - before
|
||||
if self._raise_on_retry:
|
||||
raise RuntimeError("simulated retry-credential failure")
|
||||
if self._retry_bearer is None:
|
||||
return None
|
||||
return UpstreamCredential(
|
||||
bearer=self._retry_bearer,
|
||||
base_url=self._base_url,
|
||||
expires_at="2099-01-01T00:00:00Z",
|
||||
)
|
||||
|
||||
|
||||
async def _start_runner(app: "web.Application"):
|
||||
@@ -147,6 +173,30 @@ def _build_fake_upstream(captured: Dict[str, Any]) -> "web.Application":
|
||||
return app
|
||||
|
||||
|
||||
def _build_rejecting_upstream(
|
||||
captured: Dict[str, Any], *, reject_status: int, accept_bearer: str
|
||||
) -> "web.Application":
|
||||
"""Upstream that rejects every bearer except ``accept_bearer``.
|
||||
|
||||
Drives ``handle_proxy``'s ``status in {401, 429}`` branch: the first forward
|
||||
carries the initial credential and comes back rejected, the retry carries the
|
||||
rotated one and succeeds.
|
||||
"""
|
||||
|
||||
async def gated(request):
|
||||
auth = request.headers.get("Authorization")
|
||||
captured["requests"].append({"path": request.path, "auth": auth})
|
||||
if auth != f"Bearer {accept_bearer}":
|
||||
return web.json_response(
|
||||
{"error": {"message": "rejected"}}, status=reject_status
|
||||
)
|
||||
return web.json_response({"echoed": True})
|
||||
|
||||
app = web.Application()
|
||||
app.router.add_route("*", "/v1/chat/completions", gated)
|
||||
return app
|
||||
|
||||
|
||||
async def _heartbeat(ticks: List[int], running: List[bool]) -> None:
|
||||
"""Tick a counter on the event loop until told to stop."""
|
||||
while running[0]:
|
||||
@@ -265,6 +315,144 @@ def test_credential_failure_still_maps_to_401():
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# handle_proxy -> get_retry_credential (the 401/429 rotation path)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_get_retry_credential_runs_off_the_event_loop():
|
||||
"""The 401/429 rotation must not resolve its credential on the loop thread.
|
||||
|
||||
``get_retry_credential`` is the third and last blocking method on the
|
||||
``UpstreamAdapter`` contract, and it is the most expensive of them:
|
||||
``NousPortalAdapter`` routes it into ``_get_credential(force_refresh=True)``,
|
||||
so the token-refresh POST that ``get_credential`` performs only near expiry
|
||||
is unconditional here, and it happens under the same 15s cross-process
|
||||
``_auth_store_lock()``. ``XAIGrokAdapter`` loads its key pool off disk and
|
||||
rotates it under ``self._lock``.
|
||||
|
||||
Called inline, that whole rotation runs on the loop thread and this
|
||||
assertion fails.
|
||||
"""
|
||||
async def run():
|
||||
loop_thread = threading.get_ident()
|
||||
captured: Dict[str, Any] = {"requests": []}
|
||||
upstream_runner, upstream_base = await _start_runner(
|
||||
_build_rejecting_upstream(
|
||||
captured, reject_status=401, accept_bearer="rotated-bearer"
|
||||
)
|
||||
)
|
||||
adapter = _RecordingAdapter(
|
||||
f"{upstream_base}/v1", retry_bearer="rotated-bearer"
|
||||
)
|
||||
proxy_runner, proxy_base = await _start_runner(create_app(adapter))
|
||||
try:
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.post(
|
||||
f"{proxy_base}/v1/chat/completions", json={}
|
||||
) as resp:
|
||||
assert resp.status == 200
|
||||
await resp.read()
|
||||
|
||||
assert adapter.retry_thread is not None, "get_retry_credential was never called"
|
||||
assert adapter.retry_thread != loop_thread, (
|
||||
"get_retry_credential ran on the event-loop thread "
|
||||
f"({adapter.retry_thread}); it force-refreshes the upstream token "
|
||||
"under a cross-process auth-store lock and must be offloaded"
|
||||
)
|
||||
# The rotation itself still worked: rejected bearer, then ours.
|
||||
assert adapter.retry_status_code == 401
|
||||
assert [r["auth"] for r in captured["requests"]] == [
|
||||
"Bearer test-bearer",
|
||||
"Bearer rotated-bearer",
|
||||
]
|
||||
finally:
|
||||
await proxy_runner.cleanup()
|
||||
await upstream_runner.cleanup()
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_event_loop_keeps_running_while_the_retry_credential_resolves():
|
||||
"""A stalled 429 rotation must not starve the rest of the loop.
|
||||
|
||||
Same loop-side heartbeat as the credential test, sampled by the adapter on
|
||||
entry and exit. A 429 rotation is precisely when the proxy is busiest, so
|
||||
this is the worst moment to freeze every other in-flight completion.
|
||||
"""
|
||||
async def run():
|
||||
ticks = [0]
|
||||
running = [True]
|
||||
captured: Dict[str, Any] = {"requests": []}
|
||||
upstream_runner, upstream_base = await _start_runner(
|
||||
_build_rejecting_upstream(
|
||||
captured, reject_status=429, accept_bearer="rotated-bearer"
|
||||
)
|
||||
)
|
||||
adapter = _RecordingAdapter(
|
||||
f"{upstream_base}/v1",
|
||||
stall=_STALL_SECONDS,
|
||||
ticks=ticks,
|
||||
retry_bearer="rotated-bearer",
|
||||
)
|
||||
proxy_runner, proxy_base = await _start_runner(create_app(adapter))
|
||||
beat = asyncio.create_task(_heartbeat(ticks, running))
|
||||
try:
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.post(
|
||||
f"{proxy_base}/v1/chat/completions", json={}
|
||||
) as resp:
|
||||
await resp.read()
|
||||
|
||||
assert adapter.ticks_across_retry is not None
|
||||
assert adapter.ticks_across_retry >= _MIN_TICKS_ACROSS_STALL, (
|
||||
f"only {adapter.ticks_across_retry} loop iterations ran during a "
|
||||
f"{_STALL_SECONDS}s 429 credential rotation — the event loop was frozen"
|
||||
)
|
||||
finally:
|
||||
running[0] = False
|
||||
beat.cancel()
|
||||
await asyncio.gather(beat, return_exceptions=True)
|
||||
await proxy_runner.cleanup()
|
||||
await upstream_runner.cleanup()
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_retry_credential_failure_still_returns_the_upstream_rejection():
|
||||
"""Offloading must not change the rotation's error contract.
|
||||
|
||||
``asyncio.to_thread`` re-raises the worker's exception in the awaiting
|
||||
frame, so the handler's ``except Exception -> retry_cred = None`` still
|
||||
swallows it and streams the upstream's own 401 back. Deliberately *not* in
|
||||
the red-before set — it guards behaviour the fix must leave alone.
|
||||
"""
|
||||
async def run():
|
||||
captured: Dict[str, Any] = {"requests": []}
|
||||
upstream_runner, upstream_base = await _start_runner(
|
||||
_build_rejecting_upstream(
|
||||
captured, reject_status=401, accept_bearer="never-offered"
|
||||
)
|
||||
)
|
||||
adapter = _RecordingAdapter(f"{upstream_base}/v1", raise_on_retry=True)
|
||||
proxy_runner, proxy_base = await _start_runner(create_app(adapter))
|
||||
try:
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.post(
|
||||
f"{proxy_base}/v1/chat/completions", json={}
|
||||
) as resp:
|
||||
assert resp.status == 401
|
||||
await resp.read()
|
||||
|
||||
# One forward only — the failed rotation must not be retried.
|
||||
assert len(captured["requests"]) == 1
|
||||
finally:
|
||||
await proxy_runner.cleanup()
|
||||
await upstream_runner.cleanup()
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# handle_health -> is_authenticated
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user