From db0bd42119ae41e33b120aaecae89955f0898c24 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Tue, 4 Aug 2026 13:07:41 +0530 Subject: [PATCH] fix(credential-pool): re-select in acquire_lease after a deferred refresh MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit select() re-selects once deferred single-use-token refreshes complete; acquire_lease() performed the refresh but returned its pre-refresh answer. Since _acquire_lease_under_lock returns early exactly when a refresh is pending (if not available: return None, pending_refresh), a pool whose entries all needed a refresh always returned None — the caller failed an answerable request right after the refresh succeeded. Retry once, only when the first pass was empty and a refresh ran. Post-merge gate-sweep finding on the #71775 salvage (#77714). --- agent/credential_pool.py | 8 ++ ..._credential_pool_lease_refresh_reselect.py | 115 ++++++++++++++++++ 2 files changed, 123 insertions(+) create mode 100644 tests/agent/test_credential_pool_lease_refresh_reselect.py diff --git a/agent/credential_pool.py b/agent/credential_pool.py index 1022665087..f403febf4a 100644 --- a/agent/credential_pool.py +++ b/agent/credential_pool.py @@ -2084,6 +2084,14 @@ class CredentialPool: chosen_id, pending_refresh = self._acquire_lease_under_lock(credential_id) if pending_refresh: self._refresh_pending_entries(pending_refresh) + # Mirror select(): if nothing was leasable but we just refreshed + # deferred single-use-token entries, retry now that they are back + # in rotation. Without this, a pool whose only entries all needed + # a refresh returns None even though the refresh succeeded — the + # caller sees "no credentials available" and fails a request that + # should have gone through. + if chosen_id is None: + chosen_id, _ = self._acquire_lease_under_lock(credential_id) return chosen_id def _acquire_lease_under_lock( diff --git a/tests/agent/test_credential_pool_lease_refresh_reselect.py b/tests/agent/test_credential_pool_lease_refresh_reselect.py new file mode 100644 index 0000000000..267acf9ae5 --- /dev/null +++ b/tests/agent/test_credential_pool_lease_refresh_reselect.py @@ -0,0 +1,115 @@ +"""acquire_lease must re-select after a deferred single-use-token refresh. + +Post-merge gate-sweep finding on the #71775 salvage (deferred refresh moved +OUTSIDE the pool lock). ``select()`` re-selects once the refreshed entries +are back in rotation (credential_pool.py, select() -> "if pending_refresh: +re-select"); ``acquire_lease()`` did not, so a pool whose only entries all +needed a refresh returned None even though the refresh had just succeeded — +the caller saw "no credentials available" and failed a request that should +have gone through. + +These tests stub ``_available_entries`` / ``_refresh_pending_entries`` at the +same seam the production deferred-refresh contract uses: _available_entries +returns ``(available, pending_refresh)`` and entries pending a refresh are +NOT in ``available`` until the refresh has run. +""" + +import threading + +from agent.credential_pool import CredentialPool, PooledCredential + + +def _entry(entry_id: str) -> PooledCredential: + return PooledCredential( + id=entry_id, + provider="anthropic", + auth_type="oauth", + access_token="tok", + label=entry_id, + source="oauth", + priority=0, + ) + + +def _bare_pool(entries): + """Minimal pool shell — avoids disk/keyring I/O in __init__.""" + pool = CredentialPool.__new__(CredentialPool) + pool._lock = threading.RLock() + pool._entries = list(entries) + pool._active_leases = {} + pool._current_id = None + pool._max_concurrent = 2 + pool._unmatched_rotation_streak = 0 + pool.provider = "anthropic" + return pool + + +def _wire_deferred_refresh(pool, *, refresh_succeeds: bool = True): + """Model the deferred-refresh contract with an explicit state flag.""" + state = {"needs_refresh": True, "refresh_calls": 0} + + def fake_refresh(pending): + state["refresh_calls"] += 1 + if refresh_succeeds: + state["needs_refresh"] = False + + def fake_available(clear_expired=False, refresh=False): + if state["needs_refresh"]: + # Pending a refresh -> not yet available. + pending = [(e.id, "tok") for e in pool._entries] if refresh else [] + return [], pending + return list(pool._entries), [] + + pool._refresh_pending_entries = fake_refresh + pool._available_entries = fake_available + return state + + +def test_acquire_lease_reselects_after_deferred_refresh(): + """The only entry needs a refresh; once refreshed it is available, so a + lease MUST be granted rather than reporting no credentials.""" + pool = _bare_pool([_entry("e1")]) + state = _wire_deferred_refresh(pool) + + lease = pool.acquire_lease() + + assert state["refresh_calls"] == 1, "the deferred refresh should run once" + assert state["needs_refresh"] is False, "entry is available post-refresh" + assert lease == "e1", ( + "acquire_lease returned None despite a successfully refreshed, " + "available entry — the caller would fail an answerable request" + ) + assert pool._active_leases.get("e1") == 1, "the lease must be recorded" + + +def test_acquire_lease_without_pending_refresh_does_not_double_select(): + """No pending refresh -> exactly one selection pass (no wasted work).""" + pool = _bare_pool([_entry("e1")]) + state = _wire_deferred_refresh(pool) + state["needs_refresh"] = False # already healthy + + passes = {"n": 0} + original = pool._acquire_lease_under_lock + + def counting(credential_id): + passes["n"] += 1 + return original(credential_id) + + pool._acquire_lease_under_lock = counting + + lease = pool.acquire_lease() + + assert lease == "e1" + assert passes["n"] == 1, "healthy pool must not trigger the retry path" + assert state["refresh_calls"] == 0 + + +def test_acquire_lease_still_none_when_refresh_does_not_help(): + """If the refresh leaves nothing available, None is still the answer — + the retry must not loop or invent a credential.""" + pool = _bare_pool([_entry("e1")]) + state = _wire_deferred_refresh(pool, refresh_succeeds=False) + + assert pool.acquire_lease() is None + assert state["refresh_calls"] == 1, "retry must not refresh repeatedly" + assert pool._active_leases == {}